-
Notifications
You must be signed in to change notification settings - Fork 109
/
sync_tx_receipts.go
155 lines (138 loc) · 4.48 KB
/
sync_tx_receipts.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
// Copyright Fuzamei Corp. 2018 All Rights Reserved.
// Use of this source code is governed by a BSD-style
// license that can be found in the LICENSE file.
//Package sync ...
package sync
import (
"compress/gzip"
"fmt"
"io/ioutil"
"net"
"net/http"
"strings"
"github.com/33cn/chain33/blockchain"
dbm "github.com/33cn/chain33/common/db"
l "github.com/33cn/chain33/common/log/log15"
"github.com/33cn/chain33/rpc/jsonclient"
"github.com/33cn/chain33/types"
relayerTypes "github.com/33cn/plugin/plugin/dapp/x2ethereum/ebrelayer/types"
"github.com/rs/cors"
)
var (
log = l.New("module", "sync.tx_receipts")
syncTxReceipts *TxReceipts
)
//StartSyncTxReceipt ...
func StartSyncTxReceipt(cfg *relayerTypes.SyncTxReceiptConfig, db dbm.DB) *TxReceipts {
log.Debug("StartSyncTxReceipt, load config", "para:", cfg)
log.Debug("TxReceipts started ")
bindOrResumePush(cfg)
syncTxReceipts = NewSyncTxReceipts(db)
go syncTxReceipts.SaveAndSyncTxs2Relayer()
go startHTTPService(cfg.PushBind, "*")
return syncTxReceipts
}
//func StopSyncTxReceipt() {
// syncTxReceipts.Stop()
//}
func startHTTPService(url string, clientHost string) {
listen, err := net.Listen("tcp", url)
if err != nil {
panic(err)
}
var handler http.Handler = http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
//fmt.Println(r.URL, r.Header, r.Body)
beg := types.Now()
defer func() {
log.Info("handler", "cost", types.Since(beg))
}()
client := strings.Split(r.RemoteAddr, ":")[0]
if !checkClient(client, clientHost) {
log.Error("HandlerFunc", "client", r.RemoteAddr, "expect", clientHost)
_, _ = w.Write([]byte(`{"errcode":"-1","result":null,"msg":"reject"}`))
// unbind 逻辑有问题, 需要的外部处理
// 切换外部服务时, 可能换 name
// 收到一个不是client 的请求,很有可能是以前注册过的, 取消掉
//unbind(client)
return
}
if r.URL.Path == "/" {
w.Header().Set("Content-type", "application/json")
w.WriteHeader(200)
if len(r.Header["Content-Encoding"]) >= 1 && r.Header["Content-Encoding"][0] == "gzip" {
gr, err := gzip.NewReader(r.Body)
if err != nil {
log.Debug("Error while reader serving JSON request: %v", err)
return
}
//body := make([]byte, r.ContentLength)
body, err := ioutil.ReadAll(gr)
//n, err := r.Body.Read(body)
if err != nil {
log.Debug("Error while serving JSON request: %v", err)
return
}
err = handleRequest(body)
if err == nil {
_, _ = w.Write([]byte("OK"))
} else {
_, _ = w.Write([]byte(err.Error()))
}
}
}
})
co := cors.New(cors.Options{})
handler = co.Handler(handler)
_ = http.Serve(listen, handler)
}
func handleRequest(body []byte) error {
beg := types.Now()
defer func() {
log.Info("handleRequest", "cost", types.Since(beg))
}()
var req types.TxReceipts4Subscribe
err := types.Decode(body, &req)
if err != nil {
log.Error("handleRequest", "DecodeBlockSeqErr", err)
return err
}
err = pushTxReceipts(&req)
return err
}
func checkClient(addr string, expectClient string) bool {
if expectClient == "0.0.0.0" || expectClient == "*" {
return true
}
return addr == expectClient
}
//向chain33节点的注册推送交易回执,AddSubscribeTxReceipt具有2种功能:
//首次注册功能,如果没有进行过注册,则进行首次注册
//如果已经注册,则继续推送
func bindOrResumePush(cfg *relayerTypes.SyncTxReceiptConfig) {
contract := make(map[string]bool)
contract["x2ethereum"] = true
params := types.PushSubscribeReq{
Name: cfg.PushName,
URL: cfg.PushHost,
Encode: "proto",
LastSequence: cfg.StartSyncSequence,
LastHeight: cfg.StartSyncHeight,
LastBlockHash: cfg.StartSyncHash,
Type: int32(blockchain.PushTxReceipt),
Contract: contract,
}
var res types.ReplySubscribePush
ctx := jsonclient.NewRPCCtx(cfg.Chain33Host, "Chain33.AddPushSubscribe", params, &res)
_, err := ctx.RunResult()
if err != nil {
fmt.Println("Failed to AddSubscribeTxReceipt to rpc addr:", cfg.Chain33Host, "ReplySubTxReceipt", res)
panic("bindOrResumePush client failed due to:" + err.Error())
}
if !res.IsOk {
fmt.Println("Failed to AddSubscribeTxReceipt to rpc addr:", cfg.Chain33Host, "ReplySubTxReceipt", res)
panic("bindOrResumePush client failed due to:" + res.Msg)
}
log.Info("bindOrResumePush", "Succeed to AddSubscribeTxReceipt for rpc address:", cfg.Chain33Host)
fmt.Println("Succeed to AddPushSubscribe")
}