diff --git a/broker.go b/broker.go index 703abbd..09845b9 100644 --- a/broker.go +++ b/broker.go @@ -35,8 +35,10 @@ type broker struct { termMessage chan *termMessage fileMessage chan *fileMessage userMessage chan *usrMessage - cmdMessage chan []byte - httpMessage chan *httpResp + cmdResp chan []byte + cmdReq chan *commandReq + httpResp chan *httpResp + httpReq chan *httpReq fileProxy sync.Map devCertPool *x509.CertPool } @@ -53,8 +55,10 @@ func newBroker(cfg *config.Config) *broker { termMessage: make(chan *termMessage, 1000), fileMessage: make(chan *fileMessage, 1000), userMessage: make(chan *usrMessage, 1000), - cmdMessage: make(chan []byte, 1000), - httpMessage: make(chan *httpResp, 1000), + cmdResp: make(chan []byte, 1000), + cmdReq: make(chan *commandReq, 1000), + httpResp: make(chan *httpResp, 1000), + httpReq: make(chan *httpReq, 1000), } } @@ -272,10 +276,20 @@ func (br *broker) run() { log.Error().Msg("Not found sid: " + msg.sid) } - case data := <-br.cmdMessage: + case req := <-br.cmdReq: + if dev, ok := br.devices[req.devid]; ok { + dev.WriteMsg(msgTypeCmd, req.data) + } + + case data := <-br.cmdResp: handleCmdResp(data) - case resp := <-br.httpMessage: + case req := <-br.httpReq: + if dev, ok := br.devices[req.devid]; ok { + dev.WriteMsg(msgTypeHttp, req.data) + } + + case resp := <-br.httpResp: handleHttpProxyResp(resp) } } diff --git a/command.go b/command.go index b5e5218..109b4f3 100644 --- a/command.go +++ b/command.go @@ -38,6 +38,8 @@ type commandInfo struct { type commandReq struct { cancel context.CancelFunc c *gin.Context + devid string + data []byte } var commands sync.Map @@ -68,9 +70,10 @@ func handleCmdReq(br *broker, c *gin.Context) { req := &commandReq{ cancel: cancel, c: c, + devid: devid, } - dev, ok := br.devices[devid] + _, ok := br.devices[devid] if !ok { cmdErrReply(rttyCmdErrOffline, req) return @@ -108,7 +111,8 @@ func handleCmdReq(br *broker, c *gin.Context) { msg = append(msg, 0) } - dev.WriteMsg(msgTypeCmd, msg) + req.data = msg + br.cmdReq <- req waitTime := commandTimeout diff --git a/device.go b/device.go index f52ebde..4969f54 100644 --- a/device.go +++ b/device.go @@ -278,7 +278,7 @@ func (dev *device) readLoop() { return } - dev.br.cmdMessage <- b + dev.br.cmdResp <- b case msgTypeHttp: if msgLen < 18 { @@ -286,7 +286,7 @@ func (dev *device) readLoop() { return } - dev.br.httpMessage <- &httpResp{b, dev} + dev.br.httpResp <- &httpResp{b, dev} case msgTypeHeartbeat: parseHeartbeat(dev, b) diff --git a/http.go b/http.go index e796ea6..aa33ab2 100644 --- a/http.go +++ b/http.go @@ -24,6 +24,11 @@ type httpResp struct { dev client.Client } +type httpReq struct { + devid string + data []byte +} + var httpProxyCons sync.Map var httpProxySessions sync.Map @@ -72,7 +77,8 @@ type HttpProxyWriter struct { destAddr []byte srcAddr []byte hostHeaderRewrite string - dev client.Client + br *broker + devid string } func (rw *HttpProxyWriter) Write(p []byte) (n int, err error) { @@ -80,9 +86,7 @@ func (rw *HttpProxyWriter) Write(p []byte) (n int, err error) { msg = append(msg, rw.destAddr...) msg = append(msg, p...) - dev := rw.dev.(*device) - - dev.WriteMsg(msgTypeHttp, msg) + rw.br.httpReq <- &httpReq{rw.devid, msg} return len(p), nil } @@ -108,7 +112,7 @@ func doHttpProxy(brk *broker, c net.Conn) { } devid := cookie.Value - dev, ok := brk.devices[devid] + _, ok := brk.devices[devid] if !ok { return } @@ -153,7 +157,7 @@ func doHttpProxy(brk *broker, c net.Conn) { return } - hpw := &HttpProxyWriter{destAddr, srcAddr, hostHeaderRewrite, dev} + hpw := &HttpProxyWriter{destAddr, srcAddr, hostHeaderRewrite, brk, devid} req.Host = hostHeaderRewrite hpw.WriteRequest(req)