/
subscriber.go
121 lines (92 loc) · 1.91 KB
/
subscriber.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
package publisher
import (
"encoding/hex"
"log"
"github.com/gorilla/websocket"
"golang.org/x/crypto/blake2b"
)
const channelBufSize = 100
// Chat Subscriber.
type Subscriber struct {
id string
conn *websocket.Conn
writeCh chan *Response
readCh chan *Request
closeCh chan string
doneCh chan bool
}
// Create new chat Subscriber.
func NewSubscriber(conn *websocket.Conn, readCh chan *Request, closeCh chan string) *Subscriber {
//
if conn == nil {
//
panic("conn cannot be nil")
}
hash, _ := blake2b.New256([]byte(conn.RemoteAddr().String()))
id := hex.EncodeToString(hash.Sum(nil))
writeCh := make(chan *Response, channelBufSize)
doneCh := make(chan bool)
return &Subscriber{id, conn, writeCh, readCh, closeCh, doneCh}
}
func (s *Subscriber) Conn() *websocket.Conn {
//
return s.conn
}
func (s *Subscriber) Write(r *Response) {
//
s.writeCh <- r
}
func (s *Subscriber) Del() {
//
log.Println("Deleting subscriber")
s.doneCh <- true
s.conn.Close()
}
// Listen Write and Read request via chanel
func (s *Subscriber) Listen() {
//
log.Println("Liesten to subscriber")
go s.listenWrite()
s.listenRead()
}
// Listen write request via chanel
func (s *Subscriber) listenWrite() {
//
log.Println("Listening write to Subscriber")
for {
//
select {
// send message to the Subscriber
case response := <-s.writeCh:
log.Println("Send:", response)
s.conn.WriteJSON(response)
// receive done request
case <-s.doneCh:
return
}
}
}
// Listen read request via chanel
func (s *Subscriber) listenRead() {
//
log.Println("Listening read from Subscriber")
for {
//
select {
// receive done request
case <-s.doneCh:
return
// read data from websocket connection
default:
var orf OracleRequestFormat
err := s.conn.ReadJSON(&orf)
if err != nil {
//
log.Println(err)
s.closeCh <- s.id
return
}
s.readCh <- &Request{s.id, orf}
}
}
}