@@ -26,19 +26,92 @@ function unpack(buf, opts) {
2626
2727unpack . bytes_remaining = 0 ;
2828
29+ const SEND_QUEUE_CAP = 1024 ;
30+
2931function Stream ( s ) {
3032 const self = this ;
3133 events . EventEmitter . call ( self ) ;
3234 self . buf = null ;
3335
36+ const queue = [ ] ;
37+ let waitingForDrain = false ;
38+ let flushing = false ;
39+
40+ function abandonQueue ( ) {
41+ const pending = queue . splice ( 0 , queue . length ) ;
42+ waitingForDrain = false ;
43+ if ( pending . length === 0 ) {
44+ return ;
45+ }
46+ const err = new Error (
47+ 'msgpack Stream dropped ' + pending . length +
48+ ' unsent message(s) after backpressure'
49+ ) ;
50+ for ( let i = 0 ; i < pending . length ; i ++ ) {
51+ const cb = pending [ i ] . cb ;
52+ if ( cb ) {
53+ process . nextTick ( cb , err ) ;
54+ }
55+ }
56+ self . emit ( 'error' , err ) ;
57+ }
58+
59+ function onWritableDrain ( ) {
60+ if ( flushing ) {
61+ return ;
62+ }
63+ flushing = true ;
64+ try {
65+ while ( queue . length > 0 ) {
66+ const item = queue . shift ( ) ;
67+ const ok = item . cb ? s . write ( item . buf , item . cb ) : s . write ( item . buf ) ;
68+ if ( ok === false ) {
69+ waitingForDrain = true ;
70+ return ;
71+ }
72+ }
73+ waitingForDrain = false ;
74+ self . emit ( 'drain' ) ;
75+ } finally {
76+ flushing = false ;
77+ }
78+ }
79+
3480 self . send = function ( m ) {
35- const args = [ pack ( m ) ] ;
81+ const packed = pack ( m ) ;
82+ if ( waitingForDrain ) {
83+ if ( queue . length >= SEND_QUEUE_CAP ) {
84+ throw new Error (
85+ 'msgpack Stream backpressure queue full (' +
86+ SEND_QUEUE_CAP +
87+ ' pending messages)'
88+ ) ;
89+ }
90+ let cb ;
91+ if ( arguments . length > 1 &&
92+ typeof arguments [ arguments . length - 1 ] === 'function' ) {
93+ cb = arguments [ arguments . length - 1 ] ;
94+ }
95+ queue . push ( { buf : packed , cb : cb } ) ;
96+ return false ;
97+ }
98+
99+ const args = [ packed ] ;
36100 for ( let i = 1 ; i < arguments . length ; i ++ ) {
37101 args . push ( arguments [ i ] ) ;
38102 }
39- return s . write . apply ( s , args ) ;
103+ const ok = s . write . apply ( s , args ) ;
104+ if ( ok === false ) {
105+ waitingForDrain = true ;
106+ }
107+ return ok ;
40108 } ;
41109
110+ s . addListener ( 'drain' , onWritableDrain ) ;
111+ s . addListener ( 'error' , abandonQueue ) ;
112+ s . addListener ( 'close' , abandonQueue ) ;
113+ s . addListener ( 'end' , abandonQueue ) ;
114+
42115 s . addListener ( 'data' , function ( d ) {
43116 if ( self . buf ) {
44117 const b = buffer . Buffer . allocUnsafe ( self . buf . length + d . length ) ;
0 commit comments