This repository has been archived by the owner on Sep 7, 2022. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 0
/
server.js
113 lines (97 loc) · 3.12 KB
/
server.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
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
// Copyright 2018 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
var argo = require('argo');
var router = require('argo-url-router');
var urlHelper = require('argo-url-helper');
var resource = require('argo-resource');
var root = require('./root');
var TcpReceiver = require('./tcp_receiver');
var InfluxLineParser = require('./influx_line_parser');
var http = require('http');
var qs = require('querystring');
var url = require('url');
var EventEmitter = require('events').EventEmitter;
var ws = require('ws');
var wss = ws.Server;
var socketServer = new wss({noServer: true});
var topics = {};
var server = http.createServer();
var argoServer = argo()
.use(router)
.use(urlHelper())
.use(resource(root))
argoServer = argoServer.build();
server.on('request', argoServer.run);
server.on('upgrade', function(request, socket, headers) {
var singleStreamPattern = /^\/events\?.+$/;
var globalEventStream = /^\/events$/;
if(singleStreamPattern.test(request.url)) {
socketServer.handleUpgrade(request, socket, headers, function(ws) {
var queryString = url.parse(ws.upgradeReq.url).query;
var query = qs.parse(queryString);
var topic = query.topic;
console.log(topic, ' ', typeof topic);
if(topics[topic]) {
var socketArray = topics[topic];
socketArray.push(ws);
} else {
topics[topic] = [ws];
}
ws.on('close', function() {
delete topics[topic];
});
});
} else if(globalEventStream.test(request.url)) {
socketServer.handleUpgrade(request, socket, headers, function(ws) {
var topic = '*';
if(topics[topic]) {
var socketArray = topics[topic];
socketArray.push(ws);
} else {
topics[topic] = [ws];
}
ws.on('close', function() {
delete topics['*'];
});
});
}
});
var emitter = new EventEmitter();
emitter.on('data', function(data) {
data = data.toString().split('\n');
data.forEach(function(dataLine) {
parseAndSend(dataLine);
});
function parseAndSend(line) {
var parsedData = InfluxLineParser(line);
if(parsedData) {
console.log(parsedData);
var data = parsedData.data;
if(topics[data.topic]) {
var sockets = topics[data.topic];
sockets.forEach(function(socket) {
socket.send(JSON.stringify(data));
});
}
if(topics['*']) {
var sockets = topics['*'];
sockets.forEach(function(socket) {
socket.send(JSON.stringify(data));
});
}
}
}
});
TcpReceiver(emitter, process.env.TCP_PORT || 3008);
server.listen(process.env.PORT || 3002);