fix: panic: send on closed channel

Signed-off-by: Jianhui Zhao <zhaojh329@gmail.com>
This commit is contained in:
Jianhui Zhao
2021-09-22 19:56:45 +08:00
parent 45486898f8
commit cc5e8b7896
4 changed files with 38 additions and 16 deletions
+20 -6
View File
@@ -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)
}
}
+6 -2
View File
@@ -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
+2 -2
View File
@@ -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)
+10 -6
View File
@@ -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)