heyujie
2021-05-24 4885600ecc369aa2e30a65de8dd7a410f13c34df
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
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
'use strict'
 
var handleClient
var websocket = require('websocket-stream')
var WebSocketServer = require('ws').Server
var Connection = require('mqtt-connection')
var http = require('http')
 
handleClient = function (client) {
  var self = this
 
  if (!self.clients) {
    self.clients = {}
  }
 
  client.on('connect', function (packet) {
    if (packet.clientId === 'invalid') {
      client.connack({returnCode: 2})
    } else {
      client.connack({returnCode: 0})
    }
    self.clients[packet.clientId] = client
    client.subscriptions = []
  })
 
  client.on('publish', function (packet) {
    var i, k, c, s, publish
    switch (packet.qos) {
      case 0:
        break
      case 1:
        client.puback(packet)
        break
      case 2:
        client.pubrec(packet)
        break
    }
 
    for (k in self.clients) {
      c = self.clients[k]
      publish = false
 
      for (i = 0; i < c.subscriptions.length; i++) {
        s = c.subscriptions[i]
 
        if (s.test(packet.topic)) {
          publish = true
        }
      }
 
      if (publish) {
        try {
          c.publish({topic: packet.topic, payload: packet.payload})
        } catch (error) {
          delete self.clients[k]
        }
      }
    }
  })
 
  client.on('pubrel', function (packet) {
    client.pubcomp(packet)
  })
 
  client.on('pubrec', function (packet) {
    client.pubrel(packet)
  })
 
  client.on('pubcomp', function () {
    // Nothing to be done
  })
 
  client.on('subscribe', function (packet) {
    var qos
    var topic
    var reg
    var granted = []
 
    for (var i = 0; i < packet.subscriptions.length; i++) {
      qos = packet.subscriptions[i].qos
      topic = packet.subscriptions[i].topic
      reg = new RegExp(topic.replace('+', '[^/]+').replace('#', '.+') + '$')
 
      granted.push(qos)
      client.subscriptions.push(reg)
    }
 
    client.suback({messageId: packet.messageId, granted: granted})
  })
 
  client.on('unsubscribe', function (packet) {
    client.unsuback(packet)
  })
 
  client.on('pingreq', function () {
    client.pingresp()
  })
}
 
function start (startPort, done) {
  var server = http.createServer()
  var wss = new WebSocketServer({server: server})
 
  wss.on('connection', function (ws) {
    var stream, connection
 
    if (!(ws.protocol === 'mqtt' ||
          ws.protocol === 'mqttv3.1')) {
      return ws.close()
    }
 
    stream = websocket(ws)
    connection = new Connection(stream)
    handleClient.call(server, connection)
  })
  server.listen(startPort, done)
  server.on('request', function (req, res) {
    res.statusCode = 404
    res.end('Not Found')
  })
  return server
}
 
if (require.main === module) {
  start(process.env.PORT || process.env.ZUUL_PORT, function (err) {
    if (err) {
      console.error(err)
      return
    }
    console.log('tunnelled server started on port', process.env.PORT || process.env.ZUUL_PORT)
  })
}