-
Notifications
You must be signed in to change notification settings - Fork 10
/
Copy pathadapter_io.go
100 lines (83 loc) · 1.64 KB
/
adapter_io.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
package go2p
import (
"io"
"github.com/pkg/errors"
"github.com/v-braun/awaiter"
)
type adapterIO struct {
receive chan *Message
send chan *Message
awaiter awaiter.Awaiter
adapter Adapter
emitter *eventEmitter
}
func newAdapterIO(adapter Adapter) *adapterIO {
io := new(adapterIO)
io.receive = make(chan *Message)
io.send = make(chan *Message)
io.awaiter = awaiter.New()
io.adapter = adapter
io.emitter = newEventEmitter()
return io
}
func (io *adapterIO) start() {
io.awaiter.Go(func() {
for {
m, err := io.adapter.ReadMessage()
if err != nil {
io.handleError(err, "read")
return
}
select {
case io.receive <- m:
continue
case <-io.awaiter.CancelRequested():
return
}
}
})
io.awaiter.Go(func() {
for {
select {
case m := <-io.send:
err := io.adapter.WriteMessage(m)
if err != nil {
io.handleError(err, "write")
return
}
continue
case <-io.awaiter.CancelRequested():
return
}
}
})
}
func isDisconnectErr(err error) bool {
if err == DisconnectedError || err == io.EOF {
return true
}
return false
}
func (io *adapterIO) handleError(err error, src string) {
if isDisconnectErr(err) {
io.emitter.EmitAsync("disconnect")
return
}
io.emitter.EmitAsync("error", errors.Wrapf(err, "error during %s", src))
}
func (io *adapterIO) sendMsg(m *Message) error {
select {
case io.send <- m:
return nil
case <-io.awaiter.CancelRequested():
return DisconnectedError
}
}
func (io *adapterIO) receiveMsg() (*Message, error) {
select {
case m := <-io.receive:
return m, nil
case <-io.awaiter.CancelRequested():
return nil, DisconnectedError
}
}