/
node_conn_chainbound.go
108 lines (89 loc) · 2.72 KB
/
node_conn_chainbound.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
package collector
// Plug into Chainbound fiber as mempool data source (via websocket stream):
// https://fiber.chainbound.io/docs/usage/getting-started/
import (
"context"
"time"
fiber "github.com/chainbound/fiber-go"
"github.com/flashbots/mempool-dumpster/common"
"go.uber.org/zap"
)
type ChainboundNodeOpts struct {
TxC chan TxIn
Log *zap.SugaredLogger
APIKey string
URL string // optional override, default: ChainboundDefaultURL
SourceTag string // optional override, default: "Chainbound"
}
type ChainboundNodeConnection struct {
log *zap.SugaredLogger
apiKey string
url string
srcTag string
fiberC chan *fiber.Transaction
txC chan TxIn
backoffSec int
}
func NewChainboundNodeConnection(opts ChainboundNodeOpts) *ChainboundNodeConnection {
url := opts.URL
if url == "" {
url = chainboundDefaultURL
}
srcTag := opts.SourceTag
if srcTag == "" {
srcTag = common.SourceTagChainbound
}
return &ChainboundNodeConnection{
log: opts.Log.With("src", srcTag),
apiKey: opts.APIKey,
url: url,
srcTag: srcTag,
fiberC: make(chan *fiber.Transaction),
txC: opts.TxC,
backoffSec: initialBackoffSec,
}
}
func (cbc *ChainboundNodeConnection) Start() {
cbc.log.Debug("chainbound stream starting...")
cbc.fiberC = make(chan *fiber.Transaction)
go cbc.connect()
for fiberTx := range cbc.fiberC {
nativeTx := fiberTx.ToNative()
cbc.txC <- TxIn{time.Now().UTC(), nativeTx, cbc.srcTag}
}
cbc.log.Error("chainbound stream closed")
}
func (cbc *ChainboundNodeConnection) reconnect() {
backoffDuration := time.Duration(cbc.backoffSec) * time.Second
cbc.log.Infof("reconnecting to chainbound in %s sec ...", backoffDuration.String())
time.Sleep(backoffDuration)
// increase backoff timeout for next try
cbc.backoffSec *= 2
if cbc.backoffSec > maxBackoffSec {
cbc.backoffSec = maxBackoffSec
}
cbc.Start()
}
func (cbc *ChainboundNodeConnection) connect() {
cbc.log.Infow("connecting...", "uri", cbc.url)
client := fiber.NewClient(chainboundDefaultURL, cbc.apiKey)
defer client.Close()
// Connect
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
if err := client.Connect(ctx); err != nil {
cbc.log.Errorw("failed to connect to chainbound, reconnecting in a bit...", "error", err)
go cbc.reconnect()
return
}
cbc.log.Infow("connection successful", "uri", cbc.url)
cbc.backoffSec = initialBackoffSec
// First make a sink channel on which to receive the transactions
// This is a blocking call, so it needs to run in a Goroutine
err := client.SubscribeNewTxs(nil, cbc.fiberC)
if err != nil {
cbc.log.Errorw("chainbound subscription error", "error", err)
go cbc.reconnect()
return
}
}