123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463 |
- package yamux
- import (
- "bufio"
- "fmt"
- "io"
- "math"
- "net"
- "sync"
- "sync/atomic"
- "time"
- )
- type Session struct {
-
-
- remoteGoAway int32
-
-
- localGoAway int32
-
- config *Config
-
- conn io.ReadWriteCloser
-
- bufRead *bufio.Reader
-
- pings map[uint32]chan struct{}
- pingID uint32
- pingLock sync.Mutex
-
-
- nextStreamID uint32
-
- streams map[uint32]*Stream
- streamLock sync.Mutex
-
- acceptCh chan *Stream
-
-
- sendCh chan sendReady
-
- shutdown bool
- shutdownErr error
- shutdownCh chan struct{}
- shutdownLock sync.Mutex
- }
- type sendReady struct {
- Hdr []byte
- Body io.Reader
- Err chan error
- }
- func newSession(config *Config, conn io.ReadWriteCloser, client bool) *Session {
- s := &Session{
- config: config,
- conn: conn,
- bufRead: bufio.NewReader(conn),
- pings: make(map[uint32]chan struct{}),
- streams: make(map[uint32]*Stream),
- acceptCh: make(chan *Stream, config.AcceptBacklog),
- sendCh: make(chan sendReady, 64),
- shutdownCh: make(chan struct{}),
- }
- if client {
- s.nextStreamID = 1
- } else {
- s.nextStreamID = 2
- }
- go s.recv()
- go s.send()
- if config.EnableKeepAlive {
- go s.keepalive()
- }
- return s
- }
- func (s *Session) IsClosed() bool {
- select {
- case <-s.shutdownCh:
- return true
- default:
- return false
- }
- }
- func (s *Session) Open() (*Stream, error) {
- if s.IsClosed() {
- return nil, ErrSessionShutdown
- }
- if atomic.LoadInt32(&s.remoteGoAway) == 1 {
- return nil, ErrRemoteGoAway
- }
-
- s.streamLock.Lock()
- id := s.nextStreamID
- if id >= math.MaxUint32-1 {
- s.streamLock.Unlock()
- return nil, ErrStreamsExhausted
- }
- s.nextStreamID += 2
-
- stream := newStream(s, id, streamInit)
- s.streams[id] = stream
- s.streamLock.Unlock()
-
- return stream, stream.sendWindowUpdate()
- }
- func (s *Session) Accept() (net.Conn, error) {
- return s.AcceptStream()
- }
- func (s *Session) AcceptStream() (*Stream, error) {
- select {
- case stream := <-s.acceptCh:
- return stream, nil
- case <-s.shutdownCh:
- return nil, s.shutdownErr
- }
- }
- func (s *Session) Close() error {
- s.shutdownLock.Lock()
- defer s.shutdownLock.Unlock()
- if s.shutdown {
- return nil
- }
- s.shutdown = true
- if s.shutdownErr == nil {
- s.shutdownErr = ErrSessionShutdown
- }
- close(s.shutdownCh)
- s.conn.Close()
- s.streamLock.Lock()
- defer s.streamLock.Unlock()
- for _, stream := range s.streams {
- stream.forceClose()
- }
- return nil
- }
- func (s *Session) exitErr(err error) {
- s.shutdownErr = err
- s.Close()
- }
- func (s *Session) GoAway() error {
- return s.waitForSend(s.goAway(goAwayNormal), nil)
- }
- func (s *Session) goAway(reason uint32) header {
- atomic.SwapInt32(&s.localGoAway, 1)
- hdr := header(make([]byte, headerSize))
- hdr.encode(typeGoAway, 0, 0, reason)
- return hdr
- }
- func (s *Session) Ping() (time.Duration, error) {
-
- ch := make(chan struct{})
-
- s.pingLock.Lock()
- id := s.pingID
- s.pingID++
- s.pings[id] = ch
- s.pingLock.Unlock()
-
- hdr := header(make([]byte, headerSize))
- hdr.encode(typePing, flagSYN, 0, id)
- if err := s.waitForSend(hdr, nil); err != nil {
- return 0, err
- }
-
- start := time.Now()
- select {
- case <-ch:
- case <-s.shutdownCh:
- return 0, ErrSessionShutdown
- }
-
- return time.Now().Sub(start), nil
- }
- func (s *Session) keepalive() {
- for {
- select {
- case <-time.After(s.config.KeepAliveInterval):
- s.Ping()
- case <-s.shutdownCh:
- return
- }
- }
- }
- func (s *Session) waitForSend(hdr header, body io.Reader) error {
- errCh := make(chan error, 1)
- ready := sendReady{Hdr: hdr, Body: body, Err: errCh}
- select {
- case s.sendCh <- ready:
- case <-s.shutdownCh:
- return ErrSessionShutdown
- }
- select {
- case err := <-errCh:
- return err
- case <-s.shutdownCh:
- return ErrSessionShutdown
- }
- }
- func (s *Session) sendNoWait(hdr header) error {
- select {
- case s.sendCh <- sendReady{Hdr: hdr}:
- return nil
- case <-s.shutdownCh:
- return ErrSessionShutdown
- }
- }
- func (s *Session) send() {
- for !s.IsClosed() {
- select {
- case ready := <-s.sendCh:
-
- if ready.Hdr != nil {
- sent := 0
- for sent < len(ready.Hdr) {
- n, err := s.conn.Write(ready.Hdr[sent:])
- if err != nil {
- asyncSendErr(ready.Err, err)
- s.exitErr(err)
- return
- }
- sent += n
- }
- }
-
- if ready.Body != nil {
- _, err := io.Copy(s.conn, ready.Body)
- if err != nil {
- asyncSendErr(ready.Err, err)
- s.exitErr(err)
- return
- }
- }
-
- asyncSendErr(ready.Err, nil)
- case <-s.shutdownCh:
- return
- }
- }
- }
- func (s *Session) recv() {
- hdr := header(make([]byte, headerSize))
- var handler func(header) error
- for !s.IsClosed() {
-
- if _, err := io.ReadFull(s.bufRead, hdr); err != nil {
- s.exitErr(err)
- return
- }
-
- if hdr.Version() != protoVersion {
- s.exitErr(ErrInvalidVersion)
- return
- }
-
- switch hdr.MsgType() {
- case typeData:
- handler = s.handleStreamMessage
- case typeWindowUpdate:
- handler = s.handleStreamMessage
- case typeGoAway:
- handler = s.handleGoAway
- case typePing:
- handler = s.handlePing
- default:
- s.exitErr(ErrInvalidMsgType)
- return
- }
-
- if err := handler(hdr); err != nil {
- s.exitErr(err)
- return
- }
- }
- }
- func (s *Session) handleStreamMessage(hdr header) error {
-
- id := hdr.StreamID()
- flags := hdr.Flags()
- if flags&flagSYN == flagSYN {
- if err := s.incomingStream(id); err != nil {
- return err
- }
- }
-
- s.streamLock.Lock()
- stream := s.streams[id]
- s.streamLock.Unlock()
-
- if stream == nil {
- s.sendNoWait(s.goAway(goAwayProtoErr))
- return ErrMissingStream
- }
-
- if hdr.MsgType() == typeWindowUpdate {
- if err := stream.incrSendWindow(hdr, flags); err != nil {
- s.sendNoWait(s.goAway(goAwayProtoErr))
- return err
- }
- return nil
- }
-
- if err := stream.readData(hdr, flags, s.bufRead); err != nil {
- s.sendNoWait(s.goAway(goAwayProtoErr))
- return err
- }
- return nil
- }
- func (s *Session) handlePing(hdr header) error {
- flags := hdr.Flags()
- pingID := hdr.Length()
-
- if flags&flagSYN == flagSYN {
- hdr := header(make([]byte, headerSize))
- hdr.encode(typePing, flagACK, 0, pingID)
- s.sendNoWait(hdr)
- return nil
- }
-
- s.pingLock.Lock()
- ch := s.pings[pingID]
- if ch != nil {
- delete(s.pings, pingID)
- close(ch)
- }
- s.pingLock.Unlock()
- return nil
- }
- func (s *Session) handleGoAway(hdr header) error {
- code := hdr.Length()
- switch code {
- case goAwayNormal:
- atomic.SwapInt32(&s.remoteGoAway, 1)
- case goAwayProtoErr:
- return fmt.Errorf("yamux protocol error")
- case goAwayInternalErr:
- return fmt.Errorf("remote yamux internal error")
- default:
- return fmt.Errorf("unexpected go away received")
- }
- return nil
- }
- func (s *Session) incomingStream(id uint32) error {
-
- if atomic.LoadInt32(&s.localGoAway) == 1 {
- hdr := header(make([]byte, headerSize))
- hdr.encode(typeWindowUpdate, flagRST, id, 0)
- return s.waitForSend(hdr, nil)
- }
- s.streamLock.Lock()
- defer s.streamLock.Unlock()
-
- if _, ok := s.streams[id]; ok {
- s.sendNoWait(s.goAway(goAwayProtoErr))
- return ErrDuplicateStream
- }
-
- stream := newStream(s, id, streamSYNReceived)
- s.streams[id] = stream
-
- select {
- case s.acceptCh <- stream:
- return nil
- default:
-
- delete(s.streams, id)
- stream.sendHdr.encode(typeWindowUpdate, flagRST, id, 0)
- s.sendNoWait(stream.sendHdr)
- }
- return nil
- }
- func (s *Session) closeStream(id uint32, withLock bool) {
- if !withLock {
- s.streamLock.Lock()
- defer s.streamLock.Unlock()
- }
- delete(s.streams, id)
- }
|