/
jsonLinesClient.js
95 lines (75 loc) · 2.2 KB
/
jsonLinesClient.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
'use strict';
var http = require('http'),
https = require('https');
var Parser = require('newline-json').Parser,
qs = require('qs');
var errors = require('./errors'),
FilterStream = require('./FilterStream'),
isNotHeartbeat = require('./isNotHeartbeat');
var jsonLinesClient = function (options, callback) {
var filteredStream = new FilterStream(isNotHeartbeat);
var protocol = options.protocol === 'http' ? http : https;
var req;
options.query = options.query || {};
options.query._ = Date.now();
req = protocol.request({
method: 'GET',
hostname: options.host,
port: options.port,
path: options.path + '?' + qs.stringify(options.query)
}, function (res) {
var body,
parser = new Parser();
var cleanUpStreams = function () {
req.removeAllListeners();
res.removeAllListeners();
parser.removeAllListeners();
filteredStream.removeAllListeners();
// Workaround for browserify, since its stream implementation does
// currently not know the resume function. For details see:
// https://github.com/substack/http-browserify/issues/81
if (res.resume) {
res.resume();
}
parser.end();
filteredStream.end();
};
callback({
stream: filteredStream,
disconnect: function () {
if (!res || !res.socket) {
return;
}
res.socket.end();
}
});
if (res.statusCode !== 200) {
body = '';
res.on('data', function (data) {
body += data.toString();
});
res.once('end', function () {
filteredStream.emit('error', new errors.UnexpectedStatusCode(body));
cleanUpStreams();
});
return;
}
parser.once('error', function (err) {
filteredStream.emit('error', new errors.InvalidJson('Could not parse JSON.', err));
cleanUpStreams();
});
filteredStream.once('end', function () {
cleanUpStreams();
});
res.pipe(parser).pipe(filteredStream);
});
req.once('error', function (err) {
callback({
stream: filteredStream,
disconnect: function () {}
});
filteredStream.emit('error', err);
});
req.end();
};
module.exports = jsonLinesClient;