direct.go 1.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125
  1. package ray
  2. import (
  3. "io"
  4. "sync"
  5. "time"
  6. "v2ray.com/core/common/buf"
  7. )
  8. const (
  9. bufferSize = 512
  10. )
  11. // NewRay creates a new Ray for direct traffic transport.
  12. func NewRay() Ray {
  13. return &directRay{
  14. Input: NewStream(),
  15. Output: NewStream(),
  16. }
  17. }
  18. type directRay struct {
  19. Input *Stream
  20. Output *Stream
  21. }
  22. func (v *directRay) OutboundInput() InputStream {
  23. return v.Input
  24. }
  25. func (v *directRay) OutboundOutput() OutputStream {
  26. return v.Output
  27. }
  28. func (v *directRay) InboundInput() OutputStream {
  29. return v.Input
  30. }
  31. func (v *directRay) InboundOutput() InputStream {
  32. return v.Output
  33. }
  34. type Stream struct {
  35. access sync.RWMutex
  36. closed bool
  37. buffer chan *buf.Buffer
  38. }
  39. func NewStream() *Stream {
  40. return &Stream{
  41. buffer: make(chan *buf.Buffer, bufferSize),
  42. }
  43. }
  44. func (v *Stream) Read() (*buf.Buffer, error) {
  45. if v.buffer == nil {
  46. return nil, io.EOF
  47. }
  48. v.access.RLock()
  49. if v.buffer == nil {
  50. v.access.RUnlock()
  51. return nil, io.EOF
  52. }
  53. channel := v.buffer
  54. v.access.RUnlock()
  55. result, open := <-channel
  56. if !open {
  57. return nil, io.EOF
  58. }
  59. return result, nil
  60. }
  61. func (v *Stream) Write(data *buf.Buffer) error {
  62. for !v.closed {
  63. err := v.TryWriteOnce(data)
  64. if err != io.ErrNoProgress {
  65. return err
  66. }
  67. }
  68. return io.ErrClosedPipe
  69. }
  70. func (v *Stream) TryWriteOnce(data *buf.Buffer) error {
  71. v.access.RLock()
  72. defer v.access.RUnlock()
  73. if v.closed {
  74. return io.ErrClosedPipe
  75. }
  76. select {
  77. case v.buffer <- data:
  78. return nil
  79. case <-time.After(2 * time.Second):
  80. return io.ErrNoProgress
  81. }
  82. }
  83. func (v *Stream) Close() {
  84. if v.closed {
  85. return
  86. }
  87. v.access.Lock()
  88. defer v.access.Unlock()
  89. if v.closed {
  90. return
  91. }
  92. v.closed = true
  93. close(v.buffer)
  94. }
  95. func (v *Stream) Release() {
  96. if v.buffer == nil {
  97. return
  98. }
  99. v.Close()
  100. v.access.Lock()
  101. defer v.access.Unlock()
  102. if v.buffer == nil {
  103. return
  104. }
  105. for data := range v.buffer {
  106. data.Release()
  107. }
  108. v.buffer = nil
  109. }