mirror of
https://github.com/xtaci/smux.git
synced 2024-04-21 10:51:48 +00:00
optimize read event
This commit is contained in:
+7
-17
@@ -24,13 +24,12 @@ type Session struct {
|
||||
conn io.ReadWriteCloser
|
||||
|
||||
config *Config
|
||||
nextStreamID uint32 // next stream identifier
|
||||
streams map[uint32]*Stream // all streams in this session
|
||||
rdEvents map[uint32]chan struct{} // stream read notification
|
||||
nextStreamID uint32 // next stream identifier
|
||||
|
||||
tbf chan struct{} // tokenbuffer
|
||||
frameQueues map[uint32][]Frame // stream input frame queue
|
||||
streamLock sync.Mutex // locks streams/rdEvents/frameQueues
|
||||
streams map[uint32]*Stream // all streams in this session
|
||||
streamLock sync.Mutex // locks streams && frameQueues
|
||||
|
||||
die chan struct{} // flag session has died
|
||||
dieLock sync.Mutex
|
||||
@@ -47,7 +46,6 @@ func newSession(config *Config, conn io.ReadWriteCloser, client bool) *Session {
|
||||
s.config = config
|
||||
s.streams = make(map[uint32]*Stream)
|
||||
s.frameQueues = make(map[uint32][]Frame)
|
||||
s.rdEvents = make(map[uint32]chan struct{})
|
||||
s.chAccepts = make(chan *Stream, defaultAcceptBacklog)
|
||||
s.chClosedStream = make(chan uint32, defaultCloseWait)
|
||||
s.tbf = make(chan struct{}, config.MaxFrameTokens)
|
||||
@@ -72,11 +70,9 @@ func (s *Session) OpenStream() (*Stream, error) {
|
||||
}
|
||||
|
||||
sid := atomic.AddUint32(&s.nextStreamID, 2)
|
||||
chNotifyReader := make(chan struct{}, 1)
|
||||
stream := newStream(sid, s.config.MaxFrameSize, chNotifyReader, s)
|
||||
stream := newStream(sid, s.config.MaxFrameSize, s)
|
||||
|
||||
s.streamLock.Lock()
|
||||
s.rdEvents[sid] = chNotifyReader
|
||||
s.streams[sid] = stream
|
||||
s.streamLock.Unlock()
|
||||
|
||||
@@ -187,7 +183,6 @@ func (s *Session) monitor() {
|
||||
case sid := <-s.chClosedStream:
|
||||
s.streamLock.Lock()
|
||||
delete(s.streams, sid)
|
||||
delete(s.rdEvents, sid)
|
||||
ntokens := len(s.frameQueues[sid])
|
||||
delete(s.frameQueues, sid)
|
||||
s.streamLock.Unlock()
|
||||
@@ -216,9 +211,7 @@ func (s *Session) recvLoop() {
|
||||
case cmdSYN:
|
||||
s.streamLock.Lock()
|
||||
if _, ok := s.streams[f.sid]; !ok {
|
||||
chNotifyReader := make(chan struct{}, 1)
|
||||
s.streams[f.sid] = newStream(f.sid, s.config.MaxFrameSize, chNotifyReader, s)
|
||||
s.rdEvents[f.sid] = chNotifyReader
|
||||
s.streams[f.sid] = newStream(f.sid, s.config.MaxFrameSize, s)
|
||||
s.chAccepts <- s.streams[f.sid]
|
||||
} else { // stream exists, RST the peer
|
||||
s.sendFrame(newFrame(cmdRST, f.sid))
|
||||
@@ -235,12 +228,9 @@ func (s *Session) recvLoop() {
|
||||
s.tbf <- struct{}{}
|
||||
case cmdPSH:
|
||||
s.streamLock.Lock()
|
||||
if _, ok := s.streams[f.sid]; ok {
|
||||
if stream, ok := s.streams[f.sid]; ok {
|
||||
s.frameQueues[f.sid] = append(s.frameQueues[f.sid], f)
|
||||
select {
|
||||
case s.rdEvents[f.sid] <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
stream.notifyReadEvent()
|
||||
} else { // stream is absent
|
||||
s.sendFrame(newFrame(cmdRST, f.sid))
|
||||
s.tbf <- struct{}{}
|
||||
|
||||
@@ -9,21 +9,21 @@ import (
|
||||
|
||||
// Stream implements io.ReadWriteCloser
|
||||
type Stream struct {
|
||||
id uint32
|
||||
chNotifyReader chan struct{}
|
||||
sess *Session
|
||||
frameSize uint16
|
||||
rlock sync.Mutex // read lock
|
||||
buffer []byte // temporary store of remaining frame.data
|
||||
die chan struct{} // flag the stream has closed
|
||||
dieLock sync.Mutex
|
||||
id uint32
|
||||
sess *Session
|
||||
frameSize uint16
|
||||
rlock sync.Mutex // read lock
|
||||
buffer []byte // temporary store of remaining frame.data
|
||||
chReadEvent chan struct{} // notify a read event
|
||||
die chan struct{} // flag the stream has closed
|
||||
dieLock sync.Mutex
|
||||
}
|
||||
|
||||
// newStream initiates a Stream struct
|
||||
func newStream(id uint32, frameSize uint16, chNotifyReader chan struct{}, sess *Session) *Stream {
|
||||
func newStream(id uint32, frameSize uint16, sess *Session) *Stream {
|
||||
s := new(Stream)
|
||||
s.id = id
|
||||
s.chNotifyReader = chNotifyReader
|
||||
s.chReadEvent = make(chan struct{}, 1)
|
||||
s.frameSize = frameSize
|
||||
s.sess = sess
|
||||
s.die = make(chan struct{})
|
||||
@@ -58,7 +58,7 @@ READ:
|
||||
|
||||
s.rlock.Unlock()
|
||||
select {
|
||||
case <-s.chNotifyReader:
|
||||
case <-s.chReadEvent:
|
||||
goto READ
|
||||
case <-s.die:
|
||||
return 0, errors.New(errBrokenPipe)
|
||||
@@ -120,3 +120,11 @@ func (s *Stream) split(bts []byte, cmd byte, sid uint32) []Frame {
|
||||
}
|
||||
return frames
|
||||
}
|
||||
|
||||
// notify read event
|
||||
func (s *Stream) notifyReadEvent() {
|
||||
select {
|
||||
case s.chReadEvent <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user