Skip to content

Commit afecc97

Browse files
mcollinaBridgeAR
authored andcommitted
events: add EventEmitter.on to async iterate over events
Fixes: #27847 PR-URL: #27994 Reviewed-By: Benjamin Gruenbaum <benjamingr@gmail.com> Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gus Caplan <me@gus.host> Reviewed-By: Anna Henningsen <anna@addaleax.net> Reviewed-By: Rich Trott <rtrott@gmail.com>
1 parent 07e82db commit afecc97

3 files changed

Lines changed: 362 additions & 0 deletions

File tree

‎doc/api/events.md‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -886,6 +886,41 @@ Value: `Symbol.for('nodejs.rejection')`
886886

887887
See how to write a custom [rejection handler][rejection].
888888

889+
## events.on(emitter, eventName)
890+
<!-- YAML
891+
added: REPLACEME
892+
-->
893+
894+
*`emitter` {EventEmitter}
895+
*`eventName` {string|symbol} The name of the event being listened for
896+
* Returns: {AsyncIterator} that iterates `eventName` events emitted by the `emitter`
897+
898+
```js
899+
const { on, EventEmitter } =require('events');
900+
901+
(async () => {
902+
constee=newEventEmitter();
903+
904+
// Emit later on
905+
process.nextTick(() => {
906+
ee.emit('foo', 'bar');
907+
ee.emit('foo', 42);
908+
});
909+
910+
forawait (consteventofon(ee, 'foo')) {
911+
// The execution of this inner block is synchronous and it
912+
// processes one event at a time (even with await). Do not use
913+
// if concurrent execution is required.
914+
console.log(event); // prints ['bar'] [42]
915+
}
916+
})();
917+
```
918+
919+
Returns an `AsyncIterator` that iterates `eventName` events. It will throw
920+
if the `EventEmitter` emits `'error'`. It removes all listeners when
921+
exiting the loop. The `value` returned by each iteration is an array
922+
composed of the emitted event arguments.
923+
889924
[WHATWG-EventTarget]: https://dom.spec.whatwg.org/#interface-eventtarget
890925
[`--trace-warnings`]: cli.html#cli_trace_warnings
891926
[`EventEmitter.defaultMaxListeners`]: #events_eventemitter_defaultmaxlisteners

‎lib/events.js‎

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,12 +29,16 @@ const {
2929
ObjectCreate,
3030
ObjectDefineProperty,
3131
ObjectGetPrototypeOf,
32+
ObjectSetPrototypeOf,
3233
ObjectKeys,
3334
Promise,
35+
PromiseReject,
36+
PromiseResolve,
3437
ReflectApply,
3538
ReflectOwnKeys,
3639
Symbol,
3740
SymbolFor,
41+
SymbolAsyncIterator
3842
}=primordials;
3943
constkRejection=SymbolFor('nodejs.rejection');
4044

@@ -62,6 +66,7 @@ function EventEmitter(opts) {
6266
}
6367
module.exports=EventEmitter;
6468
module.exports.once=once;
69+
module.exports.on=on;
6570

6671
// Backwards-compat with node 0.10.x
6772
EventEmitter.EventEmitter=EventEmitter;
@@ -657,3 +662,102 @@ function once(emitter, name) {
657662
emitter.once(name,eventListener);
658663
});
659664
}
665+
666+
constAsyncIteratorPrototype=ObjectGetPrototypeOf(
667+
ObjectGetPrototypeOf(asyncfunction*(){}).prototype);
668+
669+
functioncreateIterResult(value,done){
670+
return{ value, done };
671+
}
672+
673+
functionon(emitter,event){
674+
constunconsumedEvents=[];
675+
constunconsumedPromises=[];
676+
leterror=null;
677+
letfinished=false;
678+
679+
constiterator=ObjectSetPrototypeOf({
680+
next(){
681+
// First, we consume all unread events
682+
constvalue=unconsumedEvents.shift();
683+
if(value){
684+
returnPromiseResolve(createIterResult(value,false));
685+
}
686+
687+
// Then we error, if an error happened
688+
// This happens one time if at all, because after 'error'
689+
// we stop listening
690+
if(error){
691+
constp=PromiseReject(error);
692+
// Only the first element errors
693+
error=null;
694+
returnp;
695+
}
696+
697+
// If the iterator is finished, resolve to done
698+
if(finished){
699+
returnPromiseResolve(createIterResult(undefined,true));
700+
}
701+
702+
// Wait until an event happens
703+
returnnewPromise(function(resolve,reject){
704+
unconsumedPromises.push({ resolve, reject });
705+
});
706+
},
707+
708+
return(){
709+
emitter.removeListener(event,eventHandler);
710+
emitter.removeListener('error',errorHandler);
711+
finished=true;
712+
713+
for(constpromiseofunconsumedPromises){
714+
promise.resolve(createIterResult(undefined,true));
715+
}
716+
717+
returnPromiseResolve(createIterResult(undefined,true));
718+
},
719+
720+
throw(err){
721+
if(!err||!(errinstanceofError)){
722+
thrownewERR_INVALID_ARG_TYPE('EventEmitter.AsyncIterator',
723+
'Error',err);
724+
}
725+
error=err;
726+
emitter.removeListener(event,eventHandler);
727+
emitter.removeListener('error',errorHandler);
728+
},
729+
730+
[SymbolAsyncIterator](){
731+
returnthis;
732+
}
733+
},AsyncIteratorPrototype);
734+
735+
emitter.on(event,eventHandler);
736+
emitter.on('error',errorHandler);
737+
738+
returniterator;
739+
740+
functioneventHandler(...args){
741+
constpromise=unconsumedPromises.shift();
742+
if(promise){
743+
promise.resolve(createIterResult(args,false));
744+
}else{
745+
unconsumedEvents.push(args);
746+
}
747+
}
748+
749+
functionerrorHandler(err){
750+
finished=true;
751+
752+
consttoError=unconsumedPromises.shift();
753+
754+
if(toError){
755+
toError.reject(err);
756+
}else{
757+
// The next time we call next()
758+
error=err;
759+
}
760+
761+
iterator.return();
762+
}
763+
}
Lines changed: 223 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,223 @@
1+
'use strict';
2+
3+
constcommon=require('../common');
4+
constassert=require('assert');
5+
const{ on, EventEmitter }=require('events');
6+
7+
asyncfunctionbasic(){
8+
constee=newEventEmitter();
9+
process.nextTick(()=>{
10+
ee.emit('foo','bar');
11+
// 'bar' is a spurious event, we are testing
12+
// that it does not show up in the iterable
13+
ee.emit('bar',24);
14+
ee.emit('foo',42);
15+
});
16+
17+
constiterable=on(ee,'foo');
18+
19+
constexpected=[['bar'],[42]];
20+
21+
forawait(consteventofiterable){
22+
constcurrent=expected.shift();
23+
24+
assert.deepStrictEqual(current,event);
25+
26+
if(expected.length===0){
27+
break;
28+
}
29+
}
30+
assert.strictEqual(ee.listenerCount('foo'),0);
31+
assert.strictEqual(ee.listenerCount('error'),0);
32+
}
33+
34+
asyncfunctionerror(){
35+
constee=newEventEmitter();
36+
const_err=newError('kaboom');
37+
process.nextTick(()=>{
38+
ee.emit('error',_err);
39+
});
40+
41+
constiterable=on(ee,'foo');
42+
letlooped=false;
43+
letthrown=false;
44+
45+
try{
46+
// eslint-disable-next-line no-unused-vars
47+
forawait(consteventofiterable){
48+
looped=true;
49+
}
50+
}catch(err){
51+
thrown=true;
52+
assert.strictEqual(err,_err);
53+
}
54+
assert.strictEqual(thrown,true);
55+
assert.strictEqual(looped,false);
56+
}
57+
58+
asyncfunctionerrorDelayed(){
59+
constee=newEventEmitter();
60+
const_err=newError('kaboom');
61+
process.nextTick(()=>{
62+
ee.emit('foo',42);
63+
ee.emit('error',_err);
64+
});
65+
66+
constiterable=on(ee,'foo');
67+
constexpected=[[42]];
68+
letthrown=false;
69+
70+
try{
71+
forawait(consteventofiterable){
72+
constcurrent=expected.shift();
73+
assert.deepStrictEqual(current,event);
74+
}
75+
}catch(err){
76+
thrown=true;
77+
assert.strictEqual(err,_err);
78+
}
79+
assert.strictEqual(thrown,true);
80+
assert.strictEqual(ee.listenerCount('foo'),0);
81+
assert.strictEqual(ee.listenerCount('error'),0);
82+
}
83+
84+
asyncfunctionthrowInLoop(){
85+
constee=newEventEmitter();
86+
const_err=newError('kaboom');
87+
88+
process.nextTick(()=>{
89+
ee.emit('foo',42);
90+
});
91+
92+
try{
93+
forawait(consteventofon(ee,'foo')){
94+
assert.deepStrictEqual(event,[42]);
95+
throw_err;
96+
}
97+
}catch(err){
98+
assert.strictEqual(err,_err);
99+
}
100+
101+
assert.strictEqual(ee.listenerCount('foo'),0);
102+
assert.strictEqual(ee.listenerCount('error'),0);
103+
}
104+
105+
asyncfunctionnext(){
106+
constee=newEventEmitter();
107+
constiterable=on(ee,'foo');
108+
109+
process.nextTick(function(){
110+
ee.emit('foo','bar');
111+
ee.emit('foo',42);
112+
iterable.return();
113+
});
114+
115+
constresults=awaitPromise.all([
116+
iterable.next(),
117+
iterable.next(),
118+
iterable.next()
119+
]);
120+
121+
assert.deepStrictEqual(results,[{
122+
value: ['bar'],
123+
done: false
124+
},{
125+
value: [42],
126+
done: false
127+
},{
128+
value: undefined,
129+
done: true
130+
}]);
131+
132+
assert.deepStrictEqual(awaititerable.next(),{
133+
value: undefined,
134+
done: true
135+
});
136+
}
137+
138+
asyncfunctionnextError(){
139+
constee=newEventEmitter();
140+
constiterable=on(ee,'foo');
141+
const_err=newError('kaboom');
142+
process.nextTick(function(){
143+
ee.emit('error',_err);
144+
});
145+
constresults=awaitPromise.allSettled([
146+
iterable.next(),
147+
iterable.next(),
148+
iterable.next()
149+
]);
150+
assert.deepStrictEqual(results,[{
151+
status: 'rejected',
152+
reason: _err
153+
},{
154+
status: 'fulfilled',
155+
value: {
156+
value: undefined,
157+
done: true
158+
}
159+
},{
160+
status: 'fulfilled',
161+
value: {
162+
value: undefined,
163+
done: true
164+
}
165+
}]);
166+
assert.strictEqual(ee.listeners('error').length,0);
167+
}
168+
169+
asyncfunctioniterableThrow(){
170+
constee=newEventEmitter();
171+
constiterable=on(ee,'foo');
172+
173+
process.nextTick(()=>{
174+
ee.emit('foo','bar');
175+
ee.emit('foo',42);// lost in the queue
176+
iterable.throw(_err);
177+
});
178+
179+
const_err=newError('kaboom');
180+
letthrown=false;
181+
182+
assert.throws(()=>{
183+
// No argument
184+
iterable.throw();
185+
},{
186+
message: 'The "EventEmitter.AsyncIterator" property must be'+
187+
' an instance of Error. Received undefined',
188+
name: 'TypeError'
189+
});
190+
191+
constexpected=[['bar'],[42]];
192+
193+
try{
194+
forawait(consteventofiterable){
195+
assert.deepStrictEqual(event,expected.shift());
196+
}
197+
}catch(err){
198+
thrown=true;
199+
assert.strictEqual(err,_err);
200+
}
201+
assert.strictEqual(thrown,true);
202+
assert.strictEqual(expected.length,0);
203+
assert.strictEqual(ee.listenerCount('foo'),0);
204+
assert.strictEqual(ee.listenerCount('error'),0);
205+
}
206+
207+
asyncfunctionrun(){
208+
constfuncs=[
209+
basic,
210+
error,
211+
errorDelayed,
212+
throwInLoop,
213+
next,
214+
nextError,
215+
iterableThrow
216+
];
217+
218+
for(constfnoffuncs){
219+
awaitfn();
220+
}
221+
}
222+
223+
run().then(common.mustCall());

0 commit comments

Comments
 (0)