From 52ed292bb2aab5fcd695eb4428c540c28ea499a7 Mon Sep 17 00:00:00 2001 From: xtaci Date: Wed, 26 Dec 2018 22:56:12 +0800 Subject: [PATCH] del all halfway bufs, for a stable stream processor, minized jitters --- sess.go | 222 ++++++++++++++++++++++++-------------------------------- 1 file changed, 94 insertions(+), 128 deletions(-) diff --git a/sess.go b/sess.go index 75cb6e6..0c541b0 100644 --- a/sess.go +++ b/sess.go @@ -39,9 +39,6 @@ const ( // accept backlog acceptBacklog = 128 - - // prerouting(to session) queue - qlen = 1024 ) const ( @@ -627,53 +624,36 @@ func (s *UDPSession) kcpInput(data []byte) { } } -func (s *UDPSession) receiver(ch chan<- []byte) { - for { - data := xmitBuf.Get().([]byte)[:mtuLimit] - if n, _, err := s.conn.ReadFrom(data); err == nil && n >= s.headerSize+IKCP_OVERHEAD { - select { - case ch <- data[:n]: - case <-s.die: - return - } - } else if err != nil { - s.chErrorEvent <- err - return - } else { - atomic.AddUint64(&DefaultSnmp.InErrs, 1) - } - } -} - // the read loop for a client session func (s *UDPSession) readLoop() { - chPacket := make(chan []byte, qlen) - go s.receiver(chPacket) - + buf := make([]byte, mtuLimit) for { - select { - case data := <-chPacket: - raw := data - dataValid := false - if s.block != nil { - s.block.Decrypt(data, data) - data = data[nonceSize:] - checksum := crc32.ChecksumIEEE(data[crcSize:]) - if checksum == binary.LittleEndian.Uint32(data) { - data = data[crcSize:] + if n, _, err := s.conn.ReadFrom(buf); err == nil { + if n >= s.headerSize+IKCP_OVERHEAD { + data := buf[:n] + dataValid := false + if s.block != nil { + s.block.Decrypt(data, data) + data = data[nonceSize:] + checksum := crc32.ChecksumIEEE(data[crcSize:]) + if checksum == binary.LittleEndian.Uint32(data) { + data = data[crcSize:] + dataValid = true + } else { + atomic.AddUint64(&DefaultSnmp.InCsumErrors, 1) + } + } else if s.block == nil { dataValid = true - } else { - atomic.AddUint64(&DefaultSnmp.InCsumErrors, 1) } - } else if s.block == nil { - dataValid = true - } - if dataValid { - s.kcpInput(data) + if dataValid { + s.kcpInput(data) + } + } else { + atomic.AddUint64(&DefaultSnmp.InErrs, 1) } - xmitBuf.Put(raw) - case <-s.die: + } else { + s.chErrorEvent <- err return } } @@ -689,11 +669,12 @@ type ( conn net.PacketConn // the underlying packet connection sessions map[string]*UDPSession // all sessions accepted by this Listener - chAccepts chan *UDPSession // Listen() backlog - chSessionClosed chan net.Addr // session close queue - headerSize int // the additional header to a KCP frame - die chan struct{} // notify the listener has closed - rd atomic.Value // read deadline for Accept() + sessionLock sync.Mutex + chAccepts chan *UDPSession // Listen() backlog + chSessionClosed chan net.Addr // session close queue + headerSize int // the additional header to a KCP frame + die chan struct{} // notify the listener has closed + rd atomic.Value // read deadline for Accept() wd atomic.Value } @@ -709,93 +690,77 @@ func (l *Listener) monitor() { // a cache for session object last used var lastAddr string var lastSession *UDPSession - - chPacket := make(chan inPacket, qlen) - go l.receiver(chPacket) + buf := make([]byte, mtuLimit) for { - select { - case p := <-chPacket: - raw := p.data - data := p.data - from := p.from - dataValid := false - if l.block != nil { - l.block.Decrypt(data, data) - data = data[nonceSize:] - checksum := crc32.ChecksumIEEE(data[crcSize:]) - if checksum == binary.LittleEndian.Uint32(data) { - data = data[crcSize:] + if n, from, err := l.conn.ReadFrom(buf); err == nil { + if n >= l.headerSize+IKCP_OVERHEAD { + data := buf[:n] + dataValid := false + if l.block != nil { + l.block.Decrypt(data, data) + data = data[nonceSize:] + checksum := crc32.ChecksumIEEE(data[crcSize:]) + if checksum == binary.LittleEndian.Uint32(data) { + data = data[crcSize:] + dataValid = true + } else { + atomic.AddUint64(&DefaultSnmp.InCsumErrors, 1) + } + } else if l.block == nil { dataValid = true - } else { - atomic.AddUint64(&DefaultSnmp.InCsumErrors, 1) - } - } else if l.block == nil { - dataValid = true - } - - if dataValid { - addr := from.String() - var s *UDPSession - var ok bool - - // the packets received from an address always come in batch, - // cache the session for next packet, without querying map. - if addr == lastAddr { - s, ok = lastSession, true - } else if s, ok = l.sessions[addr]; ok { - lastSession = s - lastAddr = addr } - if !ok { // new session - if len(l.chAccepts) < cap(l.chAccepts) { // do not let the new sessions overwhelm accept queue - var conv uint32 - convValid := false - if l.fecDecoder != nil { - isfec := binary.LittleEndian.Uint16(data[4:]) - if isfec == typeData { - conv = binary.LittleEndian.Uint32(data[fecHeaderSizePlus2:]) + if dataValid { + addr := from.String() + var s *UDPSession + var ok bool + + // the packets received from an address always come in batch, + // cache the session for next packet, without querying map. + if addr == lastAddr { + s, ok = lastSession, true + } else { + l.sessionLock.Lock() + if s, ok = l.sessions[addr]; ok { + lastSession = s + lastAddr = addr + } + l.sessionLock.Unlock() + } + + if !ok { // new session + if len(l.chAccepts) < cap(l.chAccepts) { // do not let the new sessions overwhelm accept queue + var conv uint32 + convValid := false + if l.fecDecoder != nil { + isfec := binary.LittleEndian.Uint16(data[4:]) + if isfec == typeData { + conv = binary.LittleEndian.Uint32(data[fecHeaderSizePlus2:]) + convValid = true + } + } else { + conv = binary.LittleEndian.Uint32(data) convValid = true } - } else { - conv = binary.LittleEndian.Uint32(data) - convValid = true - } - if convValid { // creates a new session only if the 'conv' field in kcp is accessible - s := newUDPSession(conv, l.dataShards, l.parityShards, l, l.conn, from, l.block) - s.kcpInput(data) - l.sessions[addr] = s - l.chAccepts <- s + if convValid { // creates a new session only if the 'conv' field in kcp is accessible + s := newUDPSession(conv, l.dataShards, l.parityShards, l, l.conn, from, l.block) + s.kcpInput(data) + l.sessionLock.Lock() + l.sessions[addr] = s + l.sessionLock.Unlock() + l.chAccepts <- s + } } + } else { + s.kcpInput(data) } - } else { - s.kcpInput(data) } + } else { + atomic.AddUint64(&DefaultSnmp.InErrs, 1) } - - xmitBuf.Put(raw) - case deadlink := <-l.chSessionClosed: - delete(l.sessions, deadlink.String()) - case <-l.die: - return - } - } -} - -func (l *Listener) receiver(ch chan<- inPacket) { - for { - data := xmitBuf.Get().([]byte)[:mtuLimit] - if n, from, err := l.conn.ReadFrom(data); err == nil && n >= l.headerSize+IKCP_OVERHEAD { - select { - case ch <- inPacket{from, data[:n]}: - case <-l.die: - return - } - } else if err != nil { - return } else { - atomic.AddUint64(&DefaultSnmp.InErrs, 1) + return } } } @@ -872,13 +837,14 @@ func (l *Listener) Close() error { } // closeSession notify the listener that a session has closed -func (l *Listener) closeSession(remote net.Addr) bool { - select { - case l.chSessionClosed <- remote: +func (l *Listener) closeSession(remote net.Addr) (ret bool) { + l.sessionLock.Lock() + defer l.sessionLock.Unlock() + if _, ok := l.sessions[remote.String()]; ok { + delete(l.sessions, remote.String()) return true - case <-l.die: - return false } + return false } // Addr returns the listener's network address, The Addr returned is shared by all invocations of Addr, so do not modify it.