1 |
|
2 | var EventEmitter = require('events').EventEmitter;
|
3 | var util = require('util');
|
4 |
|
5 | var helpers = require('./helpers')
|
6 | var Parser = require('./parser');
|
7 | var request = require('request');
|
8 | var zlib = require('zlib');
|
9 |
|
10 | var STATUS_CODES_TO_ABORT_ON = require('./settings').STATUS_CODES_TO_ABORT_ON
|
11 |
|
12 | var StreamingAPIConnection = function (reqOpts, twitOptions) {
|
13 | this.reqOpts = reqOpts
|
14 | this.twitOptions = twitOptions
|
15 | EventEmitter.call(this)
|
16 | }
|
17 |
|
18 | util.inherits(StreamingAPIConnection, EventEmitter)
|
19 |
|
20 |
|
21 |
|
22 |
|
23 |
|
24 |
|
25 |
|
26 | StreamingAPIConnection.prototype._resetConnection = function () {
|
27 | if (this.request) {
|
28 |
|
29 | this.request.removeAllListeners();
|
30 | this.request.destroy();
|
31 | }
|
32 |
|
33 | if (this.response) {
|
34 |
|
35 | this.response.removeAllListeners();
|
36 | this.response.destroy();
|
37 | }
|
38 |
|
39 | if (this.parser) {
|
40 | this.parser.removeAllListeners()
|
41 | }
|
42 |
|
43 |
|
44 |
|
45 | clearTimeout(this._scheduledReconnect)
|
46 | delete this._scheduledReconnect
|
47 |
|
48 |
|
49 | this._stopStallAbortTimeout()
|
50 | }
|
51 |
|
52 |
|
53 |
|
54 |
|
55 | StreamingAPIConnection.prototype._resetRetryParams = function () {
|
56 |
|
57 | this._connectInterval = 0
|
58 |
|
59 | this._usedFirstReconnect = false
|
60 | }
|
61 |
|
62 | StreamingAPIConnection.prototype._startPersistentConnection = function () {
|
63 | var self = this;
|
64 | self._resetConnection();
|
65 | self._setupParser();
|
66 | self._resetStallAbortTimeout();
|
67 | self.request = request.post(this.reqOpts);
|
68 | self.emit('connect', self.request);
|
69 | self.request.on('response', function (response) {
|
70 |
|
71 |
|
72 | self._usedFirstReconnect = false;
|
73 |
|
74 | self._resetStallAbortTimeout();
|
75 | self.response = response
|
76 | if (STATUS_CODES_TO_ABORT_ON.indexOf(self.response.statusCode) !== -1) {
|
77 |
|
78 |
|
79 | var body = '';
|
80 | var compressedBody = '';
|
81 |
|
82 | self.response.on('data', function (chunk) {
|
83 | compressedBody += chunk.toString('utf8');
|
84 | })
|
85 |
|
86 | var gunzip = zlib.createGunzip();
|
87 | self.response.pipe(gunzip);
|
88 | gunzip.on('data', function (chunk) {
|
89 | body += chunk.toString('utf8')
|
90 | })
|
91 |
|
92 | gunzip.on('end', function () {
|
93 | try {
|
94 | body = JSON.parse(body)
|
95 | } catch (jsonDecodeError) {
|
96 |
|
97 |
|
98 | }
|
99 |
|
100 | var error = helpers.makeTwitError('Bad Twitter streaming request: ' + self.response.statusCode)
|
101 | error.statusCode = response ? response.statusCode: null;
|
102 | helpers.attachBodyInfoToError(error, body)
|
103 | self.emit('error', error);
|
104 |
|
105 | self.stop()
|
106 | body = null;
|
107 | });
|
108 | gunzip.on('error', function (err) {
|
109 |
|
110 |
|
111 | var errMsg = 'Gzip error: ' + err.message;
|
112 | var twitErr = helpers.makeTwitError(errMsg);
|
113 | twitErr.statusCode = self.response.statusCode;
|
114 | helpers.attachBodyInfoToError(twitErr, compressedBody);
|
115 | self.emit('error', twitErr);
|
116 | });
|
117 | } else if (self.response.statusCode === 420) {
|
118 |
|
119 | self._scheduleReconnect();
|
120 | } else {
|
121 |
|
122 |
|
123 | var gunzip = zlib.createGunzip();
|
124 | self.response.pipe(gunzip);
|
125 |
|
126 |
|
127 | gunzip.on('data', function (chunk) {
|
128 | self._connectInterval = 0
|
129 |
|
130 | self._resetStallAbortTimeout();
|
131 | self.parser.parse(chunk.toString('utf8'));
|
132 | });
|
133 |
|
134 | gunzip.on('close', self._onClose.bind(self))
|
135 | gunzip.on('error', function (err) {
|
136 | self.emit('error', err);
|
137 | })
|
138 | self.response.on('error', function (err) {
|
139 |
|
140 | self.emit('error', err);
|
141 | })
|
142 |
|
143 |
|
144 |
|
145 |
|
146 | self.emit('connected', self.response);
|
147 | }
|
148 | });
|
149 | self.request.on('close', self._onClose.bind(self));
|
150 | self.request.on('error', function (err) { self._scheduleReconnect.bind(self) });
|
151 | return self;
|
152 | }
|
153 |
|
154 |
|
155 |
|
156 |
|
157 |
|
158 |
|
159 | StreamingAPIConnection.prototype._onClose = function () {
|
160 | var self = this;
|
161 | self._stopStallAbortTimeout();
|
162 |
|
163 | if (self._abortedBy === 'twit-client')
|
164 | return
|
165 |
|
166 |
|
167 | self._abortedBy = 'twitter';
|
168 |
|
169 | if (self._scheduledReconnect) {
|
170 |
|
171 |
|
172 | return
|
173 | }
|
174 |
|
175 | self._scheduleReconnect();
|
176 | }
|
177 |
|
178 |
|
179 |
|
180 |
|
181 |
|
182 | StreamingAPIConnection.prototype.start = function () {
|
183 | this._resetRetryParams();
|
184 | this._startPersistentConnection();
|
185 | return this;
|
186 | }
|
187 |
|
188 |
|
189 |
|
190 |
|
191 |
|
192 | StreamingAPIConnection.prototype.stop = function () {
|
193 | this._abortedBy = 'twit-client';
|
194 |
|
195 | this._resetConnection();
|
196 | this._resetRetryParams();
|
197 | return this;
|
198 | }
|
199 |
|
200 |
|
201 |
|
202 |
|
203 |
|
204 |
|
205 | StreamingAPIConnection.prototype._resetStallAbortTimeout = function () {
|
206 | var self = this;
|
207 |
|
208 | self._stopStallAbortTimeout();
|
209 |
|
210 | self._stallAbortTimeout = setTimeout(function () {
|
211 | self._scheduleReconnect()
|
212 | }, 90000);
|
213 | return this;
|
214 | }
|
215 |
|
216 |
|
217 |
|
218 |
|
219 |
|
220 | StreamingAPIConnection.prototype._stopStallAbortTimeout = function () {
|
221 | clearTimeout(this._stallAbortTimeout);
|
222 |
|
223 | delete this._stallAbortTimeout;
|
224 | return this;
|
225 | }
|
226 |
|
227 |
|
228 |
|
229 |
|
230 |
|
231 |
|
232 |
|
233 | StreamingAPIConnection.prototype._scheduleReconnect = function () {
|
234 | var self = this;
|
235 | if (self.response && self.response.statusCode === 420) {
|
236 |
|
237 |
|
238 | if (!self._connectInterval) {
|
239 | self._connectInterval = 60000;
|
240 | } else {
|
241 | self._connectInterval *= 2;
|
242 | }
|
243 | } else if (self.response && String(self.response.statusCode).charAt(0) === '5') {
|
244 |
|
245 |
|
246 | if (!self._connectInterval) {
|
247 | self._connectInterval = 5000;
|
248 | } else if (self._connectInterval < 320000) {
|
249 | self._connectInterval *= 2;
|
250 | } else {
|
251 | self._connectInterval = 320000;
|
252 | }
|
253 | } else {
|
254 |
|
255 |
|
256 | if (!self._usedFirstReconnect) {
|
257 |
|
258 | self._connectInterval = 0;
|
259 | self._usedFirstReconnect = true;
|
260 | } else if (self._connectInterval < 16000) {
|
261 |
|
262 | self._connectInterval += 250;
|
263 | } else {
|
264 |
|
265 | self._connectInterval = 16000;
|
266 | }
|
267 | }
|
268 |
|
269 |
|
270 | self._scheduledReconnect = setTimeout(function () {
|
271 | self._startPersistentConnection();
|
272 | }, self._connectInterval);
|
273 | self.emit('reconnect', self.request, self.response, self._connectInterval);
|
274 | }
|
275 |
|
276 | StreamingAPIConnection.prototype._setupParser = function () {
|
277 | var self = this
|
278 | self.parser = new Parser()
|
279 |
|
280 |
|
281 |
|
282 | self.parser.on('element', function (msg) {
|
283 | self.emit('message', msg)
|
284 |
|
285 | if (msg.delete) { self.emit('delete', msg) }
|
286 | else if (msg.disconnect) { self._handleDisconnect(msg) }
|
287 | else if (msg.limit) { self.emit('limit', msg) }
|
288 | else if (msg.scrub_geo) { self.emit('scrub_geo', msg) }
|
289 | else if (msg.warning) { self.emit('warning', msg) }
|
290 | else if (msg.status_withheld) { self.emit('status_withheld', msg) }
|
291 | else if (msg.user_withheld) { self.emit('user_withheld', msg) }
|
292 | else if (msg.friends) { self.emit('friends', msg) }
|
293 | else if (msg.direct_message) { self.emit('direct_message', msg) }
|
294 | else if (msg.event) {
|
295 | self.emit('user_event', msg)
|
296 |
|
297 | var ev = msg.event
|
298 |
|
299 | if (ev === 'blocked') { self.emit('blocked', msg) }
|
300 | else if (ev === 'unblocked') { self.emit('unblocked', msg) }
|
301 | else if (ev === 'favorite') { self.emit('favorite', msg) }
|
302 | else if (ev === 'unfavorite') { self.emit('unfavorite', msg) }
|
303 | else if (ev === 'follow') { self.emit('follow', msg) }
|
304 | else if (ev === 'unfollow') { self.emit('unfollow', msg) }
|
305 | else if (ev === 'user_update') { self.emit('user_update', msg) }
|
306 | else if (ev === 'list_created') { self.emit('list_created', msg) }
|
307 | else if (ev === 'list_destroyed') { self.emit('list_destroyed', msg) }
|
308 | else if (ev === 'list_updated') { self.emit('list_updated', msg) }
|
309 | else if (ev === 'list_member_added') { self.emit('list_member_added', msg) }
|
310 | else if (ev === 'list_member_removed') { self.emit('list_member_removed', msg) }
|
311 | else if (ev === 'list_user_subscribed') { self.emit('list_user_subscribed', msg) }
|
312 | else if (ev === 'list_user_unsubscribed') { self.emit('list_user_unsubscribed', msg) }
|
313 | else if (ev === 'quoted_tweet') { self.emit('quoted_tweet', msg) }
|
314 | else { self.emit('unknown_user_event', msg) }
|
315 | } else { self.emit('tweet', msg) }
|
316 | })
|
317 |
|
318 | self.parser.on('error', function (err) {
|
319 | self.emit('parser-error', err)
|
320 | })
|
321 | }
|
322 |
|
323 | StreamingAPIConnection.prototype._handleDisconnect = function (twitterMsg) {
|
324 | this.emit('disconnect', twitterMsg);
|
325 | this.stop();
|
326 | }
|
327 |
|
328 | module.exports = StreamingAPIConnection
|