forked from bonedaddy/go-blocknative
/
subscription.go
73 lines (63 loc) · 1.48 KB
/
subscription.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
package client
import (
"fmt"
"time"
"github.com/gorilla/websocket"
)
var subscriptionSleep = 5 * time.Second
// Subscription represents a stream of events. Implementations
// carry a channel with which to store events returned by the subscription backend
type Subscription interface {
Events() chan interface{}
Unsubscribe()
Err() chan error
}
func eventLoop(cl *Client, sub Subscription, unsubMsg interface{}) {
e := sub.Events()
defer func() {
sub.Unsubscribe()
cl.WriteJSON(unsubMsg)
close(e)
}()
for {
select {
case <-sub.Err():
return
default:
var msg EthTxPayload
if err := cl.ReadJSON(&msg); err != nil {
if err := cl.ReadJSON(msg); err != nil {
if e, ok := err.(*websocket.CloseError); ok {
if e.Code != 1000 {
sub.Err() <- fmt.Errorf("websocket close error: %v", err)
return
}
}
break
} else {
break
}
}
e <- msg
}
time.Sleep(subscriptionSleep)
}
}
type subscription struct {
key string // address or txHash
eventChan chan interface{}
errChan chan error
}
// NewSubscription creates a carrier for tracking events
func NewSubscription(key string) *subscription {
return &subscription{key: key, eventChan: make(chan interface{}), errChan: make(chan error, 1)}
}
func (a *subscription) Events() chan interface{} {
return a.eventChan
}
func (a *subscription) Unsubscribe() {
a.errChan <- fmt.Errorf("subscription closed")
}
func (a *subscription) Err() chan error {
return a.errChan
}