2021-02-02 07:45:02 +00:00
|
|
|
package rpcbroadcaster
|
2020-10-08 11:50:10 +00:00
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"time"
|
|
|
|
|
2022-07-21 19:39:53 +00:00
|
|
|
"github.com/nspcc-dev/neo-go/pkg/rpcclient"
|
2020-10-08 11:50:10 +00:00
|
|
|
"go.uber.org/zap"
|
|
|
|
)
|
|
|
|
|
2022-04-20 18:30:09 +00:00
|
|
|
// RPCClient represent an rpc client for a single node.
|
2021-02-02 07:45:02 +00:00
|
|
|
type RPCClient struct {
|
2022-07-21 19:39:53 +00:00
|
|
|
client *rpcclient.Client
|
2020-10-08 11:50:10 +00:00
|
|
|
addr string
|
|
|
|
close chan struct{}
|
2022-07-01 20:04:54 +00:00
|
|
|
finished chan struct{}
|
2022-07-07 17:03:10 +00:00
|
|
|
responses chan []interface{}
|
2020-10-08 11:50:10 +00:00
|
|
|
log *zap.Logger
|
|
|
|
sendTimeout time.Duration
|
2021-02-02 07:45:02 +00:00
|
|
|
method SendMethod
|
2020-10-08 11:50:10 +00:00
|
|
|
}
|
|
|
|
|
2022-04-20 18:30:09 +00:00
|
|
|
// SendMethod represents an rpc method for sending data to other nodes.
|
2022-07-21 19:39:53 +00:00
|
|
|
type SendMethod func(*rpcclient.Client, []interface{}) error
|
2021-02-02 07:45:02 +00:00
|
|
|
|
2022-04-20 18:30:09 +00:00
|
|
|
// NewRPCClient returns a new rpc client for the provided address and method.
|
2022-07-07 17:03:10 +00:00
|
|
|
func (r *RPCBroadcaster) NewRPCClient(addr string, method SendMethod, timeout time.Duration, ch chan []interface{}) *RPCClient {
|
2021-02-02 07:45:02 +00:00
|
|
|
return &RPCClient{
|
2020-10-08 11:50:10 +00:00
|
|
|
addr: addr,
|
|
|
|
close: r.close,
|
2022-07-01 20:04:54 +00:00
|
|
|
finished: make(chan struct{}),
|
2020-10-08 11:50:10 +00:00
|
|
|
responses: ch,
|
2021-02-02 07:45:02 +00:00
|
|
|
log: r.Log.With(zap.String("address", addr)),
|
2020-10-08 11:50:10 +00:00
|
|
|
sendTimeout: timeout,
|
2021-02-02 07:45:02 +00:00
|
|
|
method: method,
|
2020-10-08 11:50:10 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-02-02 07:45:02 +00:00
|
|
|
func (c *RPCClient) run() {
|
2020-10-08 11:50:10 +00:00
|
|
|
// We ignore error as not every node can be available on startup.
|
2022-07-21 19:39:53 +00:00
|
|
|
c.client, _ = rpcclient.New(context.Background(), c.addr, rpcclient.Options{
|
2020-10-08 11:50:10 +00:00
|
|
|
DialTimeout: c.sendTimeout,
|
|
|
|
RequestTimeout: c.sendTimeout,
|
|
|
|
})
|
2022-07-01 20:04:54 +00:00
|
|
|
run:
|
2020-10-08 11:50:10 +00:00
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.close:
|
2022-07-01 20:04:54 +00:00
|
|
|
break run
|
2020-10-08 11:50:10 +00:00
|
|
|
case ps := <-c.responses:
|
|
|
|
if c.client == nil {
|
|
|
|
var err error
|
2022-07-21 19:39:53 +00:00
|
|
|
c.client, err = rpcclient.New(context.Background(), c.addr, rpcclient.Options{
|
2020-10-08 11:50:10 +00:00
|
|
|
DialTimeout: c.sendTimeout,
|
|
|
|
RequestTimeout: c.sendTimeout,
|
|
|
|
})
|
|
|
|
if err != nil {
|
2021-11-26 13:56:48 +00:00
|
|
|
c.log.Error("failed to create client to submit oracle response", zap.Error(err))
|
2020-10-08 11:50:10 +00:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
}
|
2021-02-02 07:45:02 +00:00
|
|
|
err := c.method(c.client, ps)
|
2020-10-08 11:50:10 +00:00
|
|
|
if err != nil {
|
|
|
|
c.log.Error("error while submitting oracle response", zap.Error(err))
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2022-07-01 20:07:37 +00:00
|
|
|
c.client.Close()
|
2022-07-01 20:04:54 +00:00
|
|
|
drain:
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.responses:
|
|
|
|
default:
|
|
|
|
break drain
|
|
|
|
}
|
|
|
|
}
|
|
|
|
close(c.finished)
|
2020-10-08 11:50:10 +00:00
|
|
|
}
|