-
Notifications
You must be signed in to change notification settings - Fork 0
/
direct.go
130 lines (113 loc) · 2 KB
/
direct.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
package ray
import (
"errors"
"io"
"sync"
"time"
"v2ray.com/core/common/alloc"
)
const (
bufferSize = 128
)
var (
ErrIOTimeout = errors.New("IO Timeout")
)
// NewRay creates a new Ray for direct traffic transport.
func NewRay() Ray {
return &directRay{
Input: NewStream(),
Output: NewStream(),
}
}
type directRay struct {
Input *Stream
Output *Stream
}
func (this *directRay) OutboundInput() InputStream {
return this.Input
}
func (this *directRay) OutboundOutput() OutputStream {
return this.Output
}
func (this *directRay) InboundInput() OutputStream {
return this.Input
}
func (this *directRay) InboundOutput() InputStream {
return this.Output
}
type Stream struct {
access sync.RWMutex
closed bool
buffer chan *alloc.Buffer
}
func NewStream() *Stream {
return &Stream{
buffer: make(chan *alloc.Buffer, bufferSize),
}
}
func (this *Stream) Read() (*alloc.Buffer, error) {
if this.buffer == nil {
return nil, io.EOF
}
this.access.RLock()
if this.buffer == nil {
this.access.RUnlock()
return nil, io.EOF
}
channel := this.buffer
this.access.RUnlock()
result, open := <-channel
if !open {
return nil, io.EOF
}
return result, nil
}
func (this *Stream) Write(data *alloc.Buffer) error {
for !this.closed {
err := this.TryWriteOnce(data)
if err != ErrIOTimeout {
return err
}
}
return io.EOF
}
func (this *Stream) TryWriteOnce(data *alloc.Buffer) error {
this.access.RLock()
defer this.access.RUnlock()
if this.closed {
return io.EOF
}
select {
case this.buffer <- data:
return nil
case <-time.After(2 * time.Second):
return ErrIOTimeout
}
}
func (this *Stream) Close() {
if this.closed {
return
}
this.access.Lock()
defer this.access.Unlock()
if this.closed {
return
}
this.closed = true
close(this.buffer)
}
func (this *Stream) Release() {
if this.buffer == nil {
return
}
this.Close()
this.access.Lock()
defer this.access.Unlock()
if this.buffer == nil {
return
}
for data := range this.buffer {
data.Release()
}
this.buffer = nil
}