mux.go 3.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. // Package mux multiplexes packets on a single socket (RFC7983)
  2. package mux
  3. import (
  4. "errors"
  5. "io"
  6. "net"
  7. "sync"
  8. "github.com/pion/ice/v2"
  9. "github.com/pion/logging"
  10. "github.com/pion/transport/packetio"
  11. )
  12. // The maximum amount of data that can be buffered before returning errors.
  13. const maxBufferSize = 1000 * 1000 // 1MB
  14. // Config collects the arguments to mux.Mux construction into
  15. // a single structure
  16. type Config struct {
  17. Conn net.Conn
  18. BufferSize int
  19. LoggerFactory logging.LoggerFactory
  20. }
  21. // Mux allows multiplexing
  22. type Mux struct {
  23. lock sync.RWMutex
  24. nextConn net.Conn
  25. endpoints map[*Endpoint]MatchFunc
  26. bufferSize int
  27. closedCh chan struct{}
  28. log logging.LeveledLogger
  29. }
  30. // NewMux creates a new Mux
  31. func NewMux(config Config) *Mux {
  32. m := &Mux{
  33. nextConn: config.Conn,
  34. endpoints: make(map[*Endpoint]MatchFunc),
  35. bufferSize: config.BufferSize,
  36. closedCh: make(chan struct{}),
  37. log: config.LoggerFactory.NewLogger("mux"),
  38. }
  39. go m.readLoop()
  40. return m
  41. }
  42. // NewEndpoint creates a new Endpoint
  43. func (m *Mux) NewEndpoint(f MatchFunc) *Endpoint {
  44. e := &Endpoint{
  45. mux: m,
  46. buffer: packetio.NewBuffer(),
  47. }
  48. // Set a maximum size of the buffer in bytes.
  49. e.buffer.SetLimitSize(maxBufferSize)
  50. m.lock.Lock()
  51. m.endpoints[e] = f
  52. m.lock.Unlock()
  53. return e
  54. }
  55. // RemoveEndpoint removes an endpoint from the Mux
  56. func (m *Mux) RemoveEndpoint(e *Endpoint) {
  57. m.lock.Lock()
  58. defer m.lock.Unlock()
  59. delete(m.endpoints, e)
  60. }
  61. // Close closes the Mux and all associated Endpoints.
  62. func (m *Mux) Close() error {
  63. m.lock.Lock()
  64. for e := range m.endpoints {
  65. if err := e.close(); err != nil {
  66. m.lock.Unlock()
  67. return err
  68. }
  69. delete(m.endpoints, e)
  70. }
  71. m.lock.Unlock()
  72. err := m.nextConn.Close()
  73. if err != nil {
  74. return err
  75. }
  76. // Wait for readLoop to end
  77. <-m.closedCh
  78. return nil
  79. }
  80. func (m *Mux) readLoop() {
  81. defer func() {
  82. close(m.closedCh)
  83. }()
  84. buf := make([]byte, m.bufferSize)
  85. for {
  86. n, err := m.nextConn.Read(buf)
  87. switch {
  88. case errors.Is(err, io.EOF), errors.Is(err, ice.ErrClosed):
  89. return
  90. case errors.Is(err, io.ErrShortBuffer), errors.Is(err, packetio.ErrTimeout):
  91. m.log.Errorf("mux: failed to read from packetio.Buffer %s", err.Error())
  92. continue
  93. case err != nil:
  94. m.log.Errorf("mux: ending readLoop packetio.Buffer error %s", err.Error())
  95. return
  96. }
  97. if err = m.dispatch(buf[:n]); err != nil {
  98. m.log.Errorf("mux: ending readLoop dispatch error %s", err.Error())
  99. return
  100. }
  101. }
  102. }
  103. func (m *Mux) dispatch(buf []byte) error {
  104. var endpoint *Endpoint
  105. m.lock.Lock()
  106. for e, f := range m.endpoints {
  107. if f(buf) {
  108. endpoint = e
  109. break
  110. }
  111. }
  112. m.lock.Unlock()
  113. if endpoint == nil {
  114. if len(buf) > 0 {
  115. m.log.Warnf("Warning: mux: no endpoint for packet starting with %d", buf[0])
  116. } else {
  117. m.log.Warnf("Warning: mux: no endpoint for zero length packet")
  118. }
  119. return nil
  120. }
  121. _, err := endpoint.buffer.Write(buf)
  122. // Expected when bytes are received faster than the endpoint can process them (#2152, #2180)
  123. if errors.Is(err, packetio.ErrFull) {
  124. m.log.Infof("mux: endpoint buffer is full, dropping packet")
  125. return nil
  126. }
  127. return err
  128. }