-
Notifications
You must be signed in to change notification settings - Fork 6
/
client.go
155 lines (134 loc) · 4.09 KB
/
client.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
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
package client
import (
"errors"
"fmt"
"io"
"net/url"
"os"
"os/signal"
"time"
"github.com/ElrondNetwork/elrond-go-core/core/check"
"github.com/ElrondNetwork/elrond-go-core/core/mock"
"github.com/ElrondNetwork/elrond-go-core/data/outport"
"github.com/ElrondNetwork/elrond-go-core/data/typeConverters/uint64ByteSlice"
"github.com/ElrondNetwork/elrond-go-core/marshal"
"github.com/ElrondNetwork/elrond-go-core/websocketOutportDriver"
"github.com/ElrondNetwork/elrond-go-core/websocketOutportDriver/data"
"github.com/gorilla/websocket"
)
// WSConn defines what a sender shall do
type WSConn interface {
io.Closer
ReadMessage() (messageType int, p []byte, err error)
WriteMessage(messageType int, data []byte) error
}
var (
log = &mock.LoggerMock{}
errNilMarshaller = errors.New("nil marshaller")
uint64ByteSliceConverter = uint64ByteSlice.NewBigEndianConverter()
)
type tempClient struct {
name string
marshaller marshal.Marshalizer
chanStop chan bool
}
// NewTempClient will return a new instance of tempClient
func NewTempClient(name string, marshaller marshal.Marshalizer) (*tempClient, error) {
if check.IfNil(marshaller) {
return nil, errNilMarshaller
}
return &tempClient{
name: name,
marshaller: marshaller,
chanStop: make(chan bool),
}, nil
}
// Run will start the client on the provided port
func (tc *tempClient) Run(port int) {
interrupt := make(chan os.Signal, 1)
signal.Notify(interrupt, os.Interrupt)
urlReceiveData := url.URL{Scheme: "ws", Host: fmt.Sprintf("127.0.0.1:%d", port), Path: "/operations"}
log.Info(tc.name+" -> connecting to", "url", urlReceiveData.String())
wsConnection, _, err := websocket.DefaultDialer.Dial(urlReceiveData.String(), nil)
if err != nil {
log.Error(tc.name+" -> dial", "error", err)
}
defer func() {
err = wsConnection.Close()
log.LogIfError(err)
}()
done := make(chan struct{})
go func() {
defer close(done)
for {
_, message, err := wsConnection.ReadMessage()
if err != nil {
log.Error(tc.name+" -> error read message", "error", err)
return
}
tc.verifyPayloadAndSendAckIfNeeded(message, wsConnection)
}
}()
timer := time.NewTimer(time.Second)
defer timer.Stop()
for {
select {
case <-done:
return
case <-timer.C:
case <-interrupt:
log.Info(tc.name + " -> interrupt")
// Cleanly close the connection by sending a close message and then
// waiting (with timeout) for the server to close the connection.
err = wsConnection.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""))
if err != nil {
log.Error(tc.name+" -> write close", "error", err)
return
}
select {
case <-done:
case <-time.After(time.Second):
}
return
}
}
}
func (tc *tempClient) verifyPayloadAndSendAckIfNeeded(payload []byte, ackHandler WSConn) {
if len(payload) == 0 {
log.Error(tc.name + " -> empty payload")
return
}
payloadParser, _ := websocketOutportDriver.NewWebSocketPayloadParser(uint64ByteSliceConverter)
payloadData, err := payloadParser.ExtractPayloadData(payload)
if err != nil {
log.Error(tc.name + " -> error while extracting payload data: " + err.Error())
return
}
log.Info(tc.name+" -> processing payload",
"counter", payloadData.Counter,
"operation type", payloadData.OperationType,
"message length", len(payloadData.Payload),
"data", payloadData.Payload,
)
if payloadData.OperationType.Uint32() == data.OperationSaveBlock.Uint32() {
log.Debug(tc.name + " -> save block operation")
var argsBlock outport.ArgsSaveBlockData
err = tc.marshaller.Unmarshal(&argsBlock, payload)
if err != nil {
log.Error(tc.name+" -> cannot unmarshal block", "error", err)
} else {
log.Info(tc.name+" -> successfully unmarshalled block", "hash", argsBlock.HeaderHash)
}
}
if payloadData.WithAcknowledge {
counterBytes := uint64ByteSliceConverter.ToByteSlice(payloadData.Counter)
err = ackHandler.WriteMessage(websocket.BinaryMessage, counterBytes)
if err != nil {
log.Error(tc.name + " -> " + err.Error())
}
}
}
// Stop -
func (tc *tempClient) Stop() {
tc.chanStop <- true
}