1 | 'use strict';
|
2 |
|
3 | const { Duplex } = require('stream');
|
4 |
|
5 |
|
6 |
|
7 |
|
8 |
|
9 |
|
10 |
|
11 | function emitClose(stream) {
|
12 | stream.emit('close');
|
13 | }
|
14 |
|
15 |
|
16 |
|
17 |
|
18 |
|
19 |
|
20 | function duplexOnEnd() {
|
21 | if (!this.destroyed && this._writableState.finished) {
|
22 | this.destroy();
|
23 | }
|
24 | }
|
25 |
|
26 |
|
27 |
|
28 |
|
29 |
|
30 |
|
31 |
|
32 | function duplexOnError(err) {
|
33 | this.removeListener('error', duplexOnError);
|
34 | this.destroy();
|
35 | if (this.listenerCount('error') === 0) {
|
36 |
|
37 | this.emit('error', err);
|
38 | }
|
39 | }
|
40 |
|
41 |
|
42 |
|
43 |
|
44 |
|
45 |
|
46 |
|
47 |
|
48 |
|
49 | function createWebSocketStream(ws, options) {
|
50 | let resumeOnReceiverDrain = true;
|
51 |
|
52 | function receiverOnDrain() {
|
53 | if (resumeOnReceiverDrain) ws._socket.resume();
|
54 | }
|
55 |
|
56 | if (ws.readyState === ws.CONNECTING) {
|
57 | ws.once('open', function open() {
|
58 | ws._receiver.removeAllListeners('drain');
|
59 | ws._receiver.on('drain', receiverOnDrain);
|
60 | });
|
61 | } else {
|
62 | ws._receiver.removeAllListeners('drain');
|
63 | ws._receiver.on('drain', receiverOnDrain);
|
64 | }
|
65 |
|
66 | const duplex = new Duplex({
|
67 | ...options,
|
68 | autoDestroy: false,
|
69 | emitClose: false,
|
70 | objectMode: false,
|
71 | writableObjectMode: false
|
72 | });
|
73 |
|
74 | ws.on('message', function message(msg) {
|
75 | if (!duplex.push(msg)) {
|
76 | resumeOnReceiverDrain = false;
|
77 | ws._socket.pause();
|
78 | }
|
79 | });
|
80 |
|
81 | ws.once('error', function error(err) {
|
82 | if (duplex.destroyed) return;
|
83 |
|
84 | duplex.destroy(err);
|
85 | });
|
86 |
|
87 | ws.once('close', function close() {
|
88 | if (duplex.destroyed) return;
|
89 |
|
90 | duplex.push(null);
|
91 | });
|
92 |
|
93 | duplex._destroy = function (err, callback) {
|
94 | if (ws.readyState === ws.CLOSED) {
|
95 | callback(err);
|
96 | process.nextTick(emitClose, duplex);
|
97 | return;
|
98 | }
|
99 |
|
100 | let called = false;
|
101 |
|
102 | ws.once('error', function error(err) {
|
103 | called = true;
|
104 | callback(err);
|
105 | });
|
106 |
|
107 | ws.once('close', function close() {
|
108 | if (!called) callback(err);
|
109 | process.nextTick(emitClose, duplex);
|
110 | });
|
111 | ws.terminate();
|
112 | };
|
113 |
|
114 | duplex._final = function (callback) {
|
115 | if (ws.readyState === ws.CONNECTING) {
|
116 | ws.once('open', function open() {
|
117 | duplex._final(callback);
|
118 | });
|
119 | return;
|
120 | }
|
121 |
|
122 |
|
123 |
|
124 |
|
125 |
|
126 | if (ws._socket === null) return;
|
127 |
|
128 | if (ws._socket._writableState.finished) {
|
129 | callback();
|
130 | if (duplex._readableState.endEmitted) duplex.destroy();
|
131 | } else {
|
132 | ws._socket.once('finish', function finish() {
|
133 |
|
134 |
|
135 |
|
136 | callback();
|
137 | });
|
138 | ws.close();
|
139 | }
|
140 | };
|
141 |
|
142 | duplex._read = function () {
|
143 | if (ws.readyState === ws.OPEN && !resumeOnReceiverDrain) {
|
144 | resumeOnReceiverDrain = true;
|
145 | if (!ws._receiver._writableState.needDrain) ws._socket.resume();
|
146 | }
|
147 | };
|
148 |
|
149 | duplex._write = function (chunk, encoding, callback) {
|
150 | if (ws.readyState === ws.CONNECTING) {
|
151 | ws.once('open', function open() {
|
152 | duplex._write(chunk, encoding, callback);
|
153 | });
|
154 | return;
|
155 | }
|
156 |
|
157 | ws.send(chunk, callback);
|
158 | };
|
159 |
|
160 | duplex.on('end', duplexOnEnd);
|
161 | duplex.on('error', duplexOnError);
|
162 | return duplex;
|
163 | }
|
164 |
|
165 | module.exports = createWebSocketStream;
|