Correct closing of connections

This commit is contained in:
9seconds
2019-10-10 21:07:55 +03:00
parent d10729490d
commit 56cf90b13d
3 changed files with 16 additions and 9 deletions
+1
View File
@@ -52,6 +52,7 @@ func (c *connection) write(packet conntypes.Packet) error {
func (c *connection) shutdown() { func (c *connection) shutdown() {
c.shutdownOnce.Do(func() { c.shutdownOnce.Do(func() {
c.conn.Close()
close(c.done) close(c.done)
c.hub.channelBrokenSockets <- c.id c.hub.channelBrokenSockets <- c.id
}) })
+7 -4
View File
@@ -27,14 +27,17 @@ func directConnection(request *protocol.TelegramRequest) error {
go directPipe(telegramConn, request.ClientConn, wg, request.Logger) go directPipe(telegramConn, request.ClientConn, wg, request.Logger)
go directPipe(request.ClientConn, telegramConn, wg, request.Logger) go directPipe(request.ClientConn, telegramConn, wg, request.Logger)
<-request.Ctx.Done()
wg.Wait() wg.Wait()
return request.Ctx.Err() return nil
} }
func directPipe(dst io.Writer, src io.Reader, wg *sync.WaitGroup, logger *zap.SugaredLogger) { func directPipe(dst io.WriteCloser, src io.ReadCloser, wg *sync.WaitGroup, logger *zap.SugaredLogger) {
defer wg.Done() defer func() {
dst.Close()
src.Close()
wg.Done()
}()
buf := make([]byte, directPipeBufferSize) buf := make([]byte, directPipeBufferSize)
if _, err := io.CopyBuffer(dst, src, buf); err != nil { if _, err := io.CopyBuffer(dst, src, buf); err != nil {
+8 -5
View File
@@ -32,17 +32,20 @@ func middleConnection(request *protocol.TelegramRequest) error {
go middlePipe(telegramConn, clientConn, wg, request.Logger) go middlePipe(telegramConn, clientConn, wg, request.Logger)
go middlePipe(clientConn, telegramConn, wg, request.Logger) go middlePipe(clientConn, telegramConn, wg, request.Logger)
<-request.Ctx.Done()
wg.Wait() wg.Wait()
return request.Ctx.Err() return nil
} }
func middlePipe(dst conntypes.PacketAckWriter, func middlePipe(dst conntypes.PacketAckWriteCloser,
src conntypes.PacketAckReader, src conntypes.PacketAckReadCloser,
wg *sync.WaitGroup, wg *sync.WaitGroup,
logger *zap.SugaredLogger) { logger *zap.SugaredLogger) {
defer wg.Done() defer func() {
dst.Close()
src.Close()
wg.Done()
}()
for { for {
acks := conntypes.ConnectionAcks{} acks := conntypes.ConnectionAcks{}