bridge.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190
  1. package reverse
  2. import (
  3. "context"
  4. "time"
  5. "github.com/golang/protobuf/proto"
  6. "v2ray.com/core/common/mux"
  7. "v2ray.com/core/common/net"
  8. "v2ray.com/core/common/session"
  9. "v2ray.com/core/common/task"
  10. "v2ray.com/core/common/vio"
  11. "v2ray.com/core/features/routing"
  12. "v2ray.com/core/transport/pipe"
  13. )
  14. type Bridge struct {
  15. dispatcher routing.Dispatcher
  16. tag string
  17. domain string
  18. workers []*BridgeWorker
  19. monitorTask *task.Periodic
  20. }
  21. func NewBridge(config *BridgeConfig, dispatcher routing.Dispatcher) (*Bridge, error) {
  22. if len(config.Tag) == 0 {
  23. return nil, newError("bridge tag is empty")
  24. }
  25. if len(config.Domain) == 0 {
  26. return nil, newError("bridge domain is empty")
  27. }
  28. b := &Bridge{
  29. dispatcher: dispatcher,
  30. tag: config.Tag,
  31. domain: config.Domain,
  32. }
  33. b.monitorTask = &task.Periodic{
  34. Execute: b.monitor,
  35. Interval: time.Second * 2,
  36. }
  37. return b, nil
  38. }
  39. func (b *Bridge) cleanup() {
  40. var activeWorkers []*BridgeWorker
  41. for _, w := range b.workers {
  42. if w.IsActive() {
  43. activeWorkers = append(activeWorkers, w)
  44. }
  45. }
  46. if len(activeWorkers) != len(b.workers) {
  47. b.workers = activeWorkers
  48. }
  49. }
  50. func (b *Bridge) monitor() error {
  51. b.cleanup()
  52. var numConnections uint32
  53. var numWorker uint32
  54. for _, w := range b.workers {
  55. if w.IsActive() {
  56. numConnections += w.Connections()
  57. numWorker++
  58. }
  59. }
  60. if numWorker == 0 || numConnections/numWorker > 16 {
  61. worker, err := NewBridgeWorker(b.domain, b.tag, b.dispatcher)
  62. if err != nil {
  63. newError("failed to create bridge worker").Base(err).AtWarning().WriteToLog()
  64. return nil
  65. }
  66. b.workers = append(b.workers, worker)
  67. }
  68. return nil
  69. }
  70. func (b *Bridge) Start() error {
  71. return b.monitorTask.Start()
  72. }
  73. func (b *Bridge) Close() error {
  74. return b.monitorTask.Close()
  75. }
  76. type BridgeWorker struct {
  77. tag string
  78. worker *mux.ServerWorker
  79. dispatcher routing.Dispatcher
  80. state Control_State
  81. }
  82. func NewBridgeWorker(domain string, tag string, d routing.Dispatcher) (*BridgeWorker, error) {
  83. ctx := context.Background()
  84. ctx = session.ContextWithInbound(ctx, &session.Inbound{
  85. Tag: tag,
  86. })
  87. link, err := d.Dispatch(ctx, net.Destination{
  88. Network: net.Network_TCP,
  89. Address: net.DomainAddress(domain),
  90. Port: 0,
  91. })
  92. if err != nil {
  93. return nil, err
  94. }
  95. w := &BridgeWorker{
  96. dispatcher: d,
  97. tag: tag,
  98. }
  99. worker, err := mux.NewServerWorker(context.Background(), w, link)
  100. if err != nil {
  101. return nil, err
  102. }
  103. w.worker = worker
  104. return w, nil
  105. }
  106. func (w *BridgeWorker) Type() interface{} {
  107. return routing.DispatcherType()
  108. }
  109. func (w *BridgeWorker) Start() error {
  110. return nil
  111. }
  112. func (w *BridgeWorker) Close() error {
  113. return nil
  114. }
  115. func (w *BridgeWorker) IsActive() bool {
  116. return w.state == Control_ACTIVE && !w.worker.Closed()
  117. }
  118. func (w *BridgeWorker) Connections() uint32 {
  119. return w.worker.ActiveConnections()
  120. }
  121. func (w *BridgeWorker) handleInternalConn(link vio.Link) {
  122. go func() {
  123. reader := link.Reader
  124. for {
  125. mb, err := reader.ReadMultiBuffer()
  126. if err != nil {
  127. break
  128. }
  129. for _, b := range mb {
  130. var ctl Control
  131. if err := proto.Unmarshal(b.Bytes(), &ctl); err != nil {
  132. newError("failed to parse proto message").Base(err).WriteToLog()
  133. break
  134. }
  135. if ctl.State != w.state {
  136. w.state = ctl.State
  137. }
  138. }
  139. }
  140. }()
  141. }
  142. func (w *BridgeWorker) Dispatch(ctx context.Context, dest net.Destination) (*vio.Link, error) {
  143. if !isInternalDomain(dest) {
  144. ctx = session.ContextWithInbound(ctx, &session.Inbound{
  145. Tag: w.tag,
  146. })
  147. return w.dispatcher.Dispatch(ctx, dest)
  148. }
  149. opt := []pipe.Option{pipe.WithSizeLimit(16 * 1024)}
  150. uplinkReader, uplinkWriter := pipe.New(opt...)
  151. downlinkReader, downlinkWriter := pipe.New(opt...)
  152. w.handleInternalConn(vio.Link{
  153. Reader: downlinkReader,
  154. Writer: uplinkWriter,
  155. })
  156. return &vio.Link{
  157. Reader: uplinkReader,
  158. Writer: downlinkWriter,
  159. }, nil
  160. }