diff --git a/session.go b/session.go index 0f67e3c..181b8af 100644 --- a/session.go +++ b/session.go @@ -16,12 +16,11 @@ const ( ) var ( - ErrInvalidProtocol = errors.New("invalid protocol") - ErrConsumed = errors.New("peer consumed more than sent") - ErrGoAway = errors.New("stream id overflows, should start a new connection") - ErrTimeout = errors.New("timeout") - ErrInvalidOperation = errors.New("invalid parameters on poll") - ErrWouldBlock = errors.New("operation would block on IO") + ErrInvalidProtocol = errors.New("invalid protocol") + ErrConsumed = errors.New("peer consumed more than sent") + ErrGoAway = errors.New("stream id overflows, should start a new connection") + ErrTimeout = errors.New("timeout") + ErrWouldBlock = errors.New("operation would block on IO") ) type writeRequest struct { @@ -79,15 +78,6 @@ type Session struct { shaper chan writeRequest // a shaper for writing writes chan writeRequest - - // Edge-Triggered PollIn support - // Streams which become 'readable', will return from PollWait() - pollInEvents map[uint32]*Stream - pollOutEvents map[uint32]*Stream - pollEventsLock sync.Mutex - - // stream r/w notification - chPollNotify chan struct{} } func newSession(config *Config, conn io.ReadWriteCloser, client bool) *Session { @@ -104,9 +94,6 @@ func newSession(config *Config, conn io.ReadWriteCloser, client bool) *Session { s.chSocketReadError = make(chan struct{}) s.chSocketWriteError = make(chan struct{}) s.chProtoError = make(chan struct{}) - s.chPollNotify = make(chan struct{}, 1) - s.pollInEvents = make(map[uint32]*Stream) - s.pollOutEvents = make(map[uint32]*Stream) if client { s.nextStreamID = 1 @@ -216,66 +203,6 @@ func (s *Session) Close() error { } } -// PollWait returns streams which becomes readable -func (s *Session) PollWait(revents []*Stream, wevents []*Stream) (int, int, error) { - if len(revents) == 0 || len(wevents) == 0 { - return -1, -1, ErrInvalidOperation - } - - select { - case <-s.chPollNotify: - s.pollEventsLock.Lock() - nr := 0 - for id, stream := range s.pollInEvents { - if nr >= len(revents) { - break - } - revents[nr] = stream - nr++ - delete(s.pollInEvents, id) - } - - nw := 0 - for id, stream := range s.pollOutEvents { - if nw >= len(wevents) { - break - } - wevents[nw] = stream - nw++ - delete(s.pollOutEvents, id) - } - s.pollEventsLock.Unlock() - - return nr, nw, nil - case <-s.die: - return -1, -1, io.ErrClosedPipe - } -} - -// streams notify session readable events -func (s *Session) notifyPollIn(stream *Stream) { - s.pollEventsLock.Lock() - s.pollInEvents[stream.id] = stream - s.pollEventsLock.Unlock() - - select { - case s.chPollNotify <- struct{}{}: - default: - } -} - -// streams notify session writable eevents -func (s *Session) notifyPollOut(stream *Stream) { - s.pollEventsLock.Lock() - s.pollOutEvents[stream.id] = stream - s.pollEventsLock.Unlock() - - select { - case s.chPollNotify <- struct{}{}: - default: - } -} - // notifyBucket notifies recvLoop that bucket is available func (s *Session) notifyBucket() { select { @@ -362,11 +289,6 @@ func (s *Session) streamClosed(sid uint32) { } delete(s.streams, sid) s.streamLock.Unlock() - - // poll remove - s.pollEventsLock.Lock() - delete(s.pollInEvents, sid) - s.pollEventsLock.Unlock() } // returnTokens is called by stream to return token after read diff --git a/session_test.go b/session_test.go index aacbca9..3479570 100644 --- a/session_test.go +++ b/session_test.go @@ -141,60 +141,6 @@ func TestEcho(t *testing.T) { session.Close() } -func TestPoll(t *testing.T) { - _, stop, cli, err := setupServer(t) - if err != nil { - t.Fatal(err) - } - defer stop() - session, _ := Client(cli, nil) - stream, _ := session.OpenStream() - - const N = 100 - var received int - - tx := make([]byte, 128) - go func() { - for i := 0; i < N; i++ { - stream.Write(tx) - } - }() - - buf := make([]byte, 65536) - revents := make([]*Stream, 128) - wevents := make([]*Stream, 128) - for { - n, _, err := session.PollWait(revents, wevents) - if err != nil { - log.Fatal(err) - } - - for i := 0; i < n; i++ { - stream := revents[i] - for { - size := stream.PeekSize() - nr, err := stream.TryRead(buf) - if err == ErrWouldBlock { - break - } - - if size != nr { - t.Fatal("incorrect peak") - } - - if err != nil { - t.Fatal(err) - } - received += nr - if received == len(tx)*N { - session.Close() - return - } - } - } - } -} - func TestWriteTo(t *testing.T) { const N = 1 << 20 // server diff --git a/stream.go b/stream.go index 16e644c..0f18bd9 100644 --- a/stream.go +++ b/stream.go @@ -509,10 +509,6 @@ func (s *Stream) pushBytes(buf []byte) (written int, err error) { s.bufferLock.Lock() s.buffers = append(s.buffers, buf) s.heads = append(s.heads, buf) - // Edge-Trigger - if len(s.buffers) == 1 { - s.sess.notifyPollIn(s) - } s.bufferLock.Unlock() return } @@ -546,8 +542,6 @@ func (s *Stream) update(consumed uint32, window uint32) { case s.chUpdate <- struct{}{}: default: } - - s.sess.notifyPollOut(s) } // mark this stream has been closed in protocol