mirror of
https://github.com/tendermint/tendermint.git
synced 2026-08-16 12:16:11 +00:00
Improving service endpoint
This commit is contained in:
@@ -58,7 +58,7 @@ func TestResetValidator(t *testing.T) {
|
||||
// priv val after signing is not same as empty
|
||||
assert.NotEqual(t, privVal.LastSignState, emptyState)
|
||||
|
||||
// priv val after tryAcceptConnection is same as empty
|
||||
// priv val after AcceptNewConnection is same as empty
|
||||
privVal.Reset()
|
||||
assert.Equal(t, privVal.LastSignState, emptyState)
|
||||
}
|
||||
|
||||
@@ -17,8 +17,6 @@ func RegisterRemoteSignerMsg(cdc *amino.Codec) {
|
||||
cdc.RegisterConcrete(&SignedVoteResponse{}, "tendermint/remotesigner/SignedVoteResponse", nil)
|
||||
cdc.RegisterConcrete(&SignProposalRequest{}, "tendermint/remotesigner/SignProposalRequest", nil)
|
||||
cdc.RegisterConcrete(&SignedProposalResponse{}, "tendermint/remotesigner/SignedProposalResponse", nil)
|
||||
cdc.RegisterConcrete(&PingRequest{}, "tendermint/remotesigner/PingRequest", nil)
|
||||
cdc.RegisterConcrete(&PingResponse{}, "tendermint/remotesigner/PingResponse", nil)
|
||||
}
|
||||
|
||||
// TODO: Add ChainIDRequest
|
||||
@@ -53,11 +51,3 @@ type SignedProposalResponse struct {
|
||||
Proposal *types.Proposal
|
||||
Error *RemoteSignerError
|
||||
}
|
||||
|
||||
// PingRequest is a PrivValidatorSocket message to keep the connection alive.
|
||||
type PingRequest struct {
|
||||
}
|
||||
|
||||
// PingRequest is a PrivValidatorSocket response to keep the connection alive.
|
||||
type PingResponse struct {
|
||||
}
|
||||
|
||||
@@ -85,8 +85,11 @@ func (sc *SignerClient) GetPubKey() crypto.PubKey {
|
||||
|
||||
// SignVote implements PrivValidator.
|
||||
func (sc *SignerClient) SignVote(chainID string, vote *types.Vote) error {
|
||||
sc.endpoint.Logger.Debug("SignerClient::SignVote")
|
||||
|
||||
response, err := sc.endpoint.SendRequest(&SignVoteRequest{Vote: vote})
|
||||
if err != nil {
|
||||
sc.endpoint.Logger.Error("SignerClient::SignVote", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -121,39 +124,3 @@ func (sc *SignerClient) SignProposal(chainID string, proposal *types.Proposal) e
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func handleRequest(req RemoteSignerMsg, chainID string, privVal types.PrivValidator) (RemoteSignerMsg, error) {
|
||||
var res RemoteSignerMsg
|
||||
var err error
|
||||
|
||||
switch r := req.(type) {
|
||||
case *PubKeyRequest:
|
||||
var p crypto.PubKey
|
||||
p = privVal.GetPubKey()
|
||||
res = &PubKeyResponse{p, nil}
|
||||
|
||||
case *SignVoteRequest:
|
||||
err = privVal.SignVote(chainID, r.Vote)
|
||||
if err != nil {
|
||||
res = &SignedVoteResponse{nil, &RemoteSignerError{0, err.Error()}}
|
||||
} else {
|
||||
res = &SignedVoteResponse{r.Vote, nil}
|
||||
}
|
||||
|
||||
case *SignProposalRequest:
|
||||
err = privVal.SignProposal(chainID, r.Proposal)
|
||||
if err != nil {
|
||||
res = &SignedProposalResponse{nil, &RemoteSignerError{0, err.Error()}}
|
||||
} else {
|
||||
res = &SignedProposalResponse{r.Proposal, nil}
|
||||
}
|
||||
|
||||
case *PingRequest:
|
||||
res = &PingResponse{}
|
||||
|
||||
default:
|
||||
err = fmt.Errorf("unknown msg: %v", r)
|
||||
}
|
||||
|
||||
return res, err
|
||||
}
|
||||
|
||||
@@ -14,7 +14,7 @@ type signerTestCase struct {
|
||||
chainID string
|
||||
mockPV types.PrivValidator
|
||||
signer *SignerClient
|
||||
signerService *SignerDialerEndpoint // TODO: Replace once it is encapsulated
|
||||
signerService *SignerDialerEndpoint
|
||||
}
|
||||
|
||||
func getSignerTestCases(t *testing.T) []signerTestCase {
|
||||
@@ -125,7 +125,6 @@ func TestSignerVoteResetDeadline(t *testing.T) {
|
||||
require.NoError(t, tc.signer.SignVote(tc.chainID, have))
|
||||
assert.Equal(t, want.Signature, have.Signature)
|
||||
|
||||
// FIXME: Lots of ping errors that do not bubble up
|
||||
// TODO: Clarify what is actually being tested
|
||||
|
||||
// This would exceed the deadline if it was not extended by the previous message
|
||||
@@ -148,14 +147,16 @@ func TestSignerVoteKeepAlive(t *testing.T) {
|
||||
defer tc.signerService.OnStop()
|
||||
defer tc.signer.Close()
|
||||
|
||||
// Check that even if the client does not request a
|
||||
// signature for a long time. The service is will available
|
||||
tc.signerService.Logger.Info("TEST. Forced Wait")
|
||||
time.Sleep(testTimeoutReadWrite * 2)
|
||||
tc.signerService.Logger.Info("TEST. Forced Wait - DONE")
|
||||
|
||||
require.NoError(t, tc.mockPV.SignVote(tc.chainID, want))
|
||||
require.NoError(t, tc.signer.SignVote(tc.chainID, have))
|
||||
assert.Equal(t, want.Signature, have.Signature)
|
||||
|
||||
// FIXME: Lots of ping errors that do not bubble up
|
||||
// TODO: Clarify what is actually being tested and how it differs from TestSignerVoteResetDeadline
|
||||
assert.Equal(t, want.Signature, have.Signature)
|
||||
}()
|
||||
}
|
||||
}
|
||||
@@ -214,7 +215,7 @@ type BrokenSignerDialerEndpoint struct {
|
||||
}
|
||||
|
||||
func (ss *BrokenSignerDialerEndpoint) writeMessage(msg RemoteSignerMsg) (err error) {
|
||||
_, err = cdc.MarshalBinaryLengthPrefixedWriter(ss.conn, PingResponse{})
|
||||
_, err = cdc.MarshalBinaryLengthPrefixedWriter(ss.conn, PubKeyResponse{})
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -4,13 +4,19 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/tendermint/tendermint/crypto"
|
||||
cmn "github.com/tendermint/tendermint/libs/common"
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
"github.com/tendermint/tendermint/types"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultMaxDialRetries = 10
|
||||
)
|
||||
|
||||
// SignerServiceEndpointOption sets an optional parameter on the SignerDialerEndpoint.
|
||||
type SignerServiceEndpointOption func(*SignerDialerEndpoint)
|
||||
|
||||
@@ -20,25 +26,35 @@ func SignerServiceEndpointTimeoutReadWrite(timeout time.Duration) SignerServiceE
|
||||
return func(ss *SignerDialerEndpoint) { ss.timeoutReadWrite = timeout }
|
||||
}
|
||||
|
||||
// SignerServiceEndpointConnRetries sets the amount of attempted retries to tryAcceptConnection.
|
||||
// SignerServiceEndpointConnRetries sets the amount of attempted retries to AcceptNewConnection.
|
||||
func SignerServiceEndpointConnRetries(retries int) SignerServiceEndpointOption {
|
||||
return func(ss *SignerDialerEndpoint) { ss.connRetries = retries }
|
||||
return func(ss *SignerDialerEndpoint) { ss.maxConnRetries = retries }
|
||||
}
|
||||
|
||||
// TODO: Create a type for a signerEndpoint (common for both listener/dialer)
|
||||
// TODO: Create a common type for a signerEndpoint (common for both listener/dialer)
|
||||
// getConnection
|
||||
// AcceptNewConnection
|
||||
// read
|
||||
// write
|
||||
// close
|
||||
|
||||
// SignerDialerEndpoint dials using its dialer and responds to any
|
||||
// signature requests using its privVal.
|
||||
type SignerDialerEndpoint struct {
|
||||
cmn.BaseService
|
||||
|
||||
chainID string
|
||||
timeoutReadWrite time.Duration
|
||||
connRetries int
|
||||
privVal types.PrivValidator
|
||||
|
||||
mtx sync.Mutex
|
||||
dialer SocketDialer
|
||||
conn net.Conn
|
||||
|
||||
timeoutReadWrite time.Duration
|
||||
maxConnRetries int
|
||||
|
||||
chainID string
|
||||
privVal types.PrivValidator
|
||||
|
||||
stopCh chan struct{}
|
||||
stoppedCh chan struct{}
|
||||
}
|
||||
|
||||
// NewSignerDialerEndpoint returns a SignerDialerEndpoint that will dial using the given
|
||||
@@ -50,12 +66,14 @@ func NewSignerDialerEndpoint(
|
||||
privVal types.PrivValidator,
|
||||
dialer SocketDialer,
|
||||
) *SignerDialerEndpoint {
|
||||
|
||||
se := &SignerDialerEndpoint{
|
||||
chainID: chainID,
|
||||
timeoutReadWrite: time.Second * defaultTimeoutReadWriteSeconds,
|
||||
connRetries: defaultMaxDialRetries,
|
||||
privVal: privVal,
|
||||
dialer: dialer,
|
||||
timeoutReadWrite: defaultTimeoutReadWriteSeconds * time.Second,
|
||||
maxConnRetries: defaultMaxDialRetries,
|
||||
|
||||
chainID: chainID,
|
||||
privVal: privVal,
|
||||
}
|
||||
|
||||
se.BaseService = *cmn.NewBaseService(logger, "SignerDialerEndpoint", se)
|
||||
@@ -64,51 +82,69 @@ func NewSignerDialerEndpoint(
|
||||
|
||||
// OnStart implements cmn.Service.
|
||||
func (ss *SignerDialerEndpoint) OnStart() error {
|
||||
conn, err := ss.connect()
|
||||
if err != nil {
|
||||
ss.Logger.Error("OnStart", "err", err)
|
||||
return err
|
||||
}
|
||||
ss.Logger.Debug("SignerDialerEndpoint: OnStart")
|
||||
|
||||
ss.conn = conn
|
||||
go ss.handleConnection(conn)
|
||||
ss.stopCh = make(chan struct{})
|
||||
ss.stoppedCh = make(chan struct{})
|
||||
|
||||
go ss.serviceLoop()
|
||||
|
||||
ss.Logger.Debug("OnStart - done")
|
||||
return nil
|
||||
}
|
||||
|
||||
// OnStop implements cmn.Service.
|
||||
func (ss *SignerDialerEndpoint) OnStop() {
|
||||
if ss.conn == nil {
|
||||
return
|
||||
}
|
||||
// Trigger a stop and wait
|
||||
close(ss.stopCh)
|
||||
<-ss.stoppedCh
|
||||
|
||||
if err := ss.conn.Close(); err != nil {
|
||||
ss.Logger.Error("OnStop", "err", cmn.ErrorWrap(err, "closing listener failed"))
|
||||
ss.conn = nil
|
||||
if ss.conn != nil {
|
||||
if err := ss.conn.Close(); err != nil {
|
||||
ss.Logger.Error("OnStop", "err", cmn.ErrorWrap(err, "closing listener failed"))
|
||||
ss.conn = nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *SignerDialerEndpoint) connect() (net.Conn, error) {
|
||||
for retries := 0; retries < ss.connRetries; retries++ {
|
||||
// Don't sleep if it is the first retry.
|
||||
if retries > 0 {
|
||||
time.Sleep(ss.timeoutReadWrite)
|
||||
}
|
||||
func (ss *SignerDialerEndpoint) serviceLoop() {
|
||||
defer close(ss.stoppedCh)
|
||||
|
||||
conn, err := ss.dialer()
|
||||
if err == nil {
|
||||
return conn, nil
|
||||
}
|
||||
retries := 0
|
||||
var err error = nil
|
||||
|
||||
ss.Logger.Error("dialing", "err", err)
|
||||
for {
|
||||
select {
|
||||
default:
|
||||
{
|
||||
if retries > ss.maxConnRetries {
|
||||
ss.Logger.Error("Maximum retries reached", "retries", retries)
|
||||
return
|
||||
}
|
||||
|
||||
if ss.conn == nil {
|
||||
ss.conn, err = ss.dialer()
|
||||
if err != nil {
|
||||
ss.Logger.Error("SignerDialerEndpoint::serviceLoop", "err", err)
|
||||
retries += 1
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
retries = 0
|
||||
ss.handleRequest()
|
||||
}
|
||||
|
||||
case <-ss.stopCh:
|
||||
{
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return nil, ErrDialRetryMax
|
||||
}
|
||||
|
||||
func (ss *SignerDialerEndpoint) readMessage() (msg RemoteSignerMsg, err error) {
|
||||
// TODO: Avoid duplication. Unify endpoints
|
||||
|
||||
if ss.conn == nil {
|
||||
return nil, fmt.Errorf("not connected")
|
||||
}
|
||||
@@ -153,33 +189,64 @@ func (ss *SignerDialerEndpoint) writeMessage(msg RemoteSignerMsg) (err error) {
|
||||
return
|
||||
}
|
||||
|
||||
func (ss *SignerDialerEndpoint) handleConnection(conn net.Conn) {
|
||||
for {
|
||||
if !ss.IsRunning() {
|
||||
return // Ignore error from listener closing.
|
||||
func (ss *SignerDialerEndpoint) handleRequest() {
|
||||
if !ss.IsRunning() {
|
||||
return // Ignore error from listener closing.
|
||||
}
|
||||
|
||||
ss.Logger.Info("SignerDialerEndpoint: connected", "timeout", ss.timeoutReadWrite)
|
||||
|
||||
req, err := ss.readMessage()
|
||||
if err != nil {
|
||||
if err != io.EOF {
|
||||
ss.Logger.Error("SignerDialerEndpoint handleMessage", "err", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
ss.Logger.Debug("SignerDialerEndpoint: connected", "timeout", ss.timeoutReadWrite)
|
||||
res, err := handleMessage(req, ss.chainID, ss.privVal)
|
||||
|
||||
req, err := ss.readMessage()
|
||||
if err != nil {
|
||||
if err != io.EOF {
|
||||
ss.Logger.Error("SignerDialerEndpoint handleConnection", "err", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
// only log the error; we'll reply with an error in res
|
||||
ss.Logger.Error("handleMessage handleMessage", "err", err)
|
||||
}
|
||||
|
||||
res, err := handleRequest(req, ss.chainID, ss.privVal)
|
||||
|
||||
if err != nil {
|
||||
// only log the error; we'll reply with an error in res
|
||||
ss.Logger.Error("handleConnection handleRequest", "err", err)
|
||||
}
|
||||
|
||||
err = ss.writeMessage(res)
|
||||
if err != nil {
|
||||
ss.Logger.Error("handleConnection writeMessage", "err", err)
|
||||
return
|
||||
}
|
||||
err = ss.writeMessage(res)
|
||||
if err != nil {
|
||||
ss.Logger.Error("handleMessage writeMessage", "err", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func handleMessage(req RemoteSignerMsg, chainID string, privVal types.PrivValidator) (RemoteSignerMsg, error) {
|
||||
var res RemoteSignerMsg
|
||||
var err error
|
||||
|
||||
switch r := req.(type) {
|
||||
case *PubKeyRequest:
|
||||
var p crypto.PubKey
|
||||
p = privVal.GetPubKey()
|
||||
res = &PubKeyResponse{p, nil}
|
||||
|
||||
case *SignVoteRequest:
|
||||
err = privVal.SignVote(chainID, r.Vote)
|
||||
if err != nil {
|
||||
res = &SignedVoteResponse{nil, &RemoteSignerError{0, err.Error()}}
|
||||
} else {
|
||||
res = &SignedVoteResponse{r.Vote, nil}
|
||||
}
|
||||
|
||||
case *SignProposalRequest:
|
||||
err = privVal.SignProposal(chainID, r.Proposal)
|
||||
if err != nil {
|
||||
res = &SignedProposalResponse{nil, &RemoteSignerError{0, err.Error()}}
|
||||
} else {
|
||||
res = &SignedProposalResponse{r.Proposal, nil}
|
||||
}
|
||||
|
||||
default:
|
||||
err = fmt.Errorf("unknown msg: %v", r)
|
||||
}
|
||||
|
||||
return res, err
|
||||
}
|
||||
|
||||
@@ -10,35 +10,12 @@ import (
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultHeartbeatSeconds = 2
|
||||
defaultMaxDialRetries = 10
|
||||
)
|
||||
|
||||
var (
|
||||
heartbeatPeriod = time.Second * defaultHeartbeatSeconds
|
||||
)
|
||||
|
||||
// SignerValidatorEndpointOption sets an optional parameter on the SocketVal.
|
||||
type SignerValidatorEndpointOption func(*SignerListenerEndpoint)
|
||||
|
||||
// SignerValidatorEndpointSetHeartbeat sets the period on which to check the liveness of the
|
||||
// connected Signer connections.
|
||||
func SignerValidatorEndpointSetHeartbeat(period time.Duration) SignerValidatorEndpointOption {
|
||||
return func(sc *SignerListenerEndpoint) { sc.heartbeatPeriod = period }
|
||||
}
|
||||
|
||||
// TODO: Add a type for SignerEndpoints
|
||||
// getConnection
|
||||
// tryAcceptConnection
|
||||
// read
|
||||
// write
|
||||
// close
|
||||
|
||||
// TODO: Fix comments
|
||||
// SocketVal implements PrivValidator.
|
||||
// It listens for an external process to dial in and uses
|
||||
// the socket to request signatures.
|
||||
// SignerListenerEndpoint listens for an external process to dial in
|
||||
// and keeps the connection alive by dropping and reconnecting
|
||||
type SignerListenerEndpoint struct {
|
||||
cmn.BaseService
|
||||
|
||||
@@ -47,19 +24,14 @@ type SignerListenerEndpoint struct {
|
||||
conn net.Conn
|
||||
|
||||
timeoutReadWrite time.Duration
|
||||
|
||||
// ping
|
||||
cancelPingCh chan struct{}
|
||||
pingTicker *time.Ticker
|
||||
heartbeatPeriod time.Duration
|
||||
}
|
||||
|
||||
// NewSignerListenerEndpoint returns an instance of SignerListenerEndpoint.
|
||||
func NewSignerListenerEndpoint(logger log.Logger, listener net.Listener) *SignerListenerEndpoint {
|
||||
|
||||
sc := &SignerListenerEndpoint{
|
||||
listener: listener,
|
||||
timeoutReadWrite : defaultTimeoutReadWriteSeconds * time.Second,
|
||||
heartbeatPeriod: heartbeatPeriod,
|
||||
}
|
||||
|
||||
sc.BaseService = *cmn.NewBaseService(logger, "SignerListenerEndpoint", sc)
|
||||
@@ -71,67 +43,18 @@ func NewSignerListenerEndpoint(logger log.Logger, listener net.Listener) *Signer
|
||||
func (ve *SignerListenerEndpoint) OnStart() error {
|
||||
ve.Logger.Debug("SignerListenerEndpoint: OnStart")
|
||||
|
||||
err := ve.tryAcceptConnection()
|
||||
err := ve.AcceptNewConnection()
|
||||
if err != nil {
|
||||
ve.Logger.Error("OnStart", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
ve.Logger.Info("OnStart", "connected", ve.isConnected())
|
||||
|
||||
// Start a routine to keep the connection alive
|
||||
go ve.pingLoop()
|
||||
|
||||
ve.Logger.Debug("SignerListenerEndpoint OnStart", "connected", ve.isConnected())
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ve *SignerListenerEndpoint)pingLoop() {
|
||||
ve.cancelPingCh = make(chan struct{}, 1)
|
||||
ve.pingTicker = time.NewTicker(ve.heartbeatPeriod)
|
||||
|
||||
fmt.Printf("ONSTART out -> %p\n", ve)
|
||||
|
||||
for {
|
||||
fmt.Printf("ONSTART in -> %p\n", ve)
|
||||
|
||||
select {
|
||||
case <-ve.pingTicker.C:
|
||||
fmt.Printf("ONSTART before -> %p\n", ve)
|
||||
err := ve.ping()
|
||||
|
||||
if err != nil {
|
||||
ve.Logger.Error("Ping", "err", err)
|
||||
if err == ErrUnexpectedResponse {
|
||||
return
|
||||
}
|
||||
|
||||
err := ve.tryAcceptConnection()
|
||||
if err != nil {
|
||||
ve.Logger.Error("Connection from remote signer not available", "err", err)
|
||||
continue
|
||||
}
|
||||
|
||||
ve.Logger.Info("Connection from remote signer available", "impl", ve)
|
||||
}
|
||||
case <-ve.cancelPingCh:
|
||||
ve.Logger.Debug("SignerListenerEndpoint: cancel ping ch")
|
||||
|
||||
ve.pingTicker.Stop()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// OnStop implements cmn.Service.
|
||||
func (ve *SignerListenerEndpoint) OnStop() {
|
||||
ve.Logger.Debug("SignerListenerEndpoint: OnStop")
|
||||
|
||||
if ve.cancelPingCh != nil {
|
||||
ve.Logger.Debug("SignerListenerEndpoint: Close cancel ping channel")
|
||||
close(ve.cancelPingCh)
|
||||
ve.cancelPingCh = nil
|
||||
}
|
||||
|
||||
ve.Logger.Debug("SignerListenerEndpoint: OnStop calling Close")
|
||||
_ = ve.Close()
|
||||
}
|
||||
@@ -156,14 +79,7 @@ func (ve *SignerListenerEndpoint) Close() error {
|
||||
|
||||
ve.Logger.Debug("SignerListenerEndpoint: Close")
|
||||
|
||||
if ve.conn != nil {
|
||||
if err := ve.conn.Close(); err != nil {
|
||||
ve.Logger.Error("Closing connection", "err", err)
|
||||
return err
|
||||
}
|
||||
ve.Logger.Debug("SignerListenerEndpoint: set ve.conn Nil")
|
||||
ve.conn = nil
|
||||
}
|
||||
ve.dropConnection()
|
||||
|
||||
if ve.listener != nil {
|
||||
if err := ve.listener.Close(); err != nil {
|
||||
@@ -183,7 +99,11 @@ func (ve *SignerListenerEndpoint) SendRequest(request RemoteSignerMsg) (RemoteSi
|
||||
ve.Logger.Debug("SignerListenerEndpoint: Send request", "connected", ve.isConnected())
|
||||
|
||||
if !ve.isConnected() {
|
||||
return nil, cmn.ErrorWrap(ErrListenerNoConnection, "endpoint is not connected")
|
||||
ve.Logger.Info("SignerListenerEndpoint: Reconnecting")
|
||||
err := ve.AcceptNewConnection()
|
||||
if err != nil {
|
||||
return nil, cmn.ErrorWrap(ErrListenerNoConnection, "could not reconnect")
|
||||
}
|
||||
}
|
||||
|
||||
ve.Logger.Debug("Send request. Write")
|
||||
@@ -197,6 +117,7 @@ func (ve *SignerListenerEndpoint) SendRequest(request RemoteSignerMsg) (RemoteSi
|
||||
|
||||
res, err := ve.readMessage()
|
||||
if err != nil {
|
||||
ve.Logger.Debug("Read Error", "err", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -209,6 +130,18 @@ func (ve *SignerListenerEndpoint) isConnected() bool {
|
||||
return ve.conn != nil
|
||||
}
|
||||
|
||||
// dropConnection closes the current connection but does not touch the listening socket
|
||||
func (ve *SignerListenerEndpoint) dropConnection() {
|
||||
ve.Logger.Debug("SignerListenerEndpoint: dropConnection")
|
||||
|
||||
if ve.conn != nil {
|
||||
if err := ve.conn.Close(); err != nil {
|
||||
ve.Logger.Error("Closing connection", "err", err)
|
||||
}
|
||||
ve.conn = nil
|
||||
}
|
||||
}
|
||||
|
||||
func (ve *SignerListenerEndpoint) readMessage() (msg RemoteSignerMsg, err error) {
|
||||
if !ve.isConnected() {
|
||||
return nil, cmn.ErrorWrap(ErrListenerNoConnection, "endpoint is not connected")
|
||||
@@ -216,7 +149,11 @@ func (ve *SignerListenerEndpoint) readMessage() (msg RemoteSignerMsg, err error)
|
||||
|
||||
// Reset read deadline
|
||||
deadline := time.Now().Add(ve.timeoutReadWrite)
|
||||
ve.Logger.Debug("SignerListenerEndpoint: readMessage", "deadline", deadline)
|
||||
ve.Logger.Debug(
|
||||
"SignerListenerEndpoint: readMessage",
|
||||
"timeout", ve.timeoutReadWrite,
|
||||
"deadline", deadline)
|
||||
|
||||
err = ve.conn.SetReadDeadline(deadline)
|
||||
if err != nil {
|
||||
return
|
||||
@@ -226,6 +163,7 @@ func (ve *SignerListenerEndpoint) readMessage() (msg RemoteSignerMsg, err error)
|
||||
_, err = cdc.UnmarshalBinaryLengthPrefixedReader(ve.conn, &msg, maxRemoteSignerMsgSize)
|
||||
if _, ok := err.(timeoutError); ok {
|
||||
err = cmn.ErrorWrap(ErrListenerTimeout, err.Error())
|
||||
ve.dropConnection()
|
||||
}
|
||||
|
||||
return
|
||||
@@ -236,11 +174,13 @@ func (ve *SignerListenerEndpoint) writeMessage(msg RemoteSignerMsg) (err error)
|
||||
return cmn.ErrorWrap(ErrListenerNoConnection, "endpoint is not connected")
|
||||
}
|
||||
|
||||
fmt.Printf("writemessage -> %p\n", ve)
|
||||
|
||||
// Reset read deadline
|
||||
deadline := time.Now().Add(ve.timeoutReadWrite)
|
||||
ve.Logger.Debug("SignerListenerEndpoint: writeMessage", "deadline", deadline)
|
||||
ve.Logger.Debug(
|
||||
"SignerListenerEndpoint: writeMessage",
|
||||
"timeout", ve.timeoutReadWrite,
|
||||
"deadline", deadline)
|
||||
|
||||
err = ve.conn.SetWriteDeadline(deadline)
|
||||
if err != nil {
|
||||
return
|
||||
@@ -254,12 +194,12 @@ func (ve *SignerListenerEndpoint) writeMessage(msg RemoteSignerMsg) (err error)
|
||||
return
|
||||
}
|
||||
|
||||
// tryAcceptConnection waits to accept a new connection.
|
||||
func (ve *SignerListenerEndpoint) tryAcceptConnection() error {
|
||||
// AcceptNewConnection waits to accept a new connection.
|
||||
func (ve *SignerListenerEndpoint) AcceptNewConnection() error {
|
||||
ve.mtx.Lock()
|
||||
defer ve.mtx.Unlock()
|
||||
|
||||
ve.Logger.Debug("SignerListenerEndpoint: tryAcceptConnection")
|
||||
ve.Logger.Debug("SignerListenerEndpoint: AcceptNewConnection")
|
||||
|
||||
if !ve.IsRunning() || ve.listener == nil {
|
||||
return fmt.Errorf("endpoint is closing")
|
||||
@@ -268,13 +208,12 @@ func (ve *SignerListenerEndpoint) tryAcceptConnection() error {
|
||||
// if the conn already exists and close it.
|
||||
if ve.conn != nil {
|
||||
if tmpErr := ve.conn.Close(); tmpErr != nil {
|
||||
ve.Logger.Error("tryAcceptConnection: error closing old connection", "err", tmpErr)
|
||||
ve.Logger.Error("AcceptNewConnection: error closing old connection", "err", tmpErr)
|
||||
}
|
||||
}
|
||||
|
||||
// Forget old connection
|
||||
ve.conn = nil
|
||||
ve.Logger.Debug("SignerListenerEndpoint: set ve.conn Nil")
|
||||
|
||||
// wait for a new conn
|
||||
conn, err := ve.listener.Accept()
|
||||
@@ -284,30 +223,7 @@ func (ve *SignerListenerEndpoint) tryAcceptConnection() error {
|
||||
}
|
||||
|
||||
ve.conn = conn
|
||||
ve.Logger.Debug("SignerListenerEndpoint: New connection", "connected", ve.isConnected())
|
||||
|
||||
fmt.Printf("tryAcceptConnection: %p\n", ve)
|
||||
// TODO: maybe we need to inform the owner that a connection has been received
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Ping is used to check connection health.
|
||||
func (ve *SignerListenerEndpoint) ping() error {
|
||||
ve.Logger.Debug("SignerListenerEndpoint: PING", "connected", ve.isConnected())
|
||||
|
||||
response, err := ve.SendRequest(&PingRequest{})
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
_, ok := response.(*PingResponse)
|
||||
if !ok {
|
||||
return ErrUnexpectedResponse
|
||||
}
|
||||
|
||||
ve.Logger.Debug("SignerListenerEndpoint: pong")
|
||||
ve.Logger.Info("SignerListenerEndpoint: New connection")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -18,9 +18,6 @@ var (
|
||||
|
||||
testTimeoutReadWrite = 100 * time.Millisecond
|
||||
testTimeoutReadWrite2o3 = 66 * time.Millisecond // 2/3 of the other one
|
||||
|
||||
testTimeoutHeartbeat = 10 * time.Millisecond
|
||||
testTimeoutHeartbeat3o2 = 6 * time.Millisecond // 3/2 of the other one
|
||||
)
|
||||
|
||||
type dialerTestCase struct {
|
||||
@@ -31,7 +28,7 @@ type dialerTestCase struct {
|
||||
// TestSignerRemoteRetryTCPOnly will test connection retry attempts over TCP. We
|
||||
// don't need this for Unix sockets because the OS instantly knows the state of
|
||||
// both ends of the socket connection. This basically causes the
|
||||
// SignerDialerEndpoint.dialer() call inside SignerDialerEndpoint.tryAcceptConnection() to return
|
||||
// SignerDialerEndpoint.dialer() call inside SignerDialerEndpoint.AcceptNewConnection() to return
|
||||
// successfully immediately, putting an instant stop to any retry attempts.
|
||||
func TestSignerRemoteRetryTCPOnly(t *testing.T) {
|
||||
var (
|
||||
@@ -98,8 +95,6 @@ func TestRetryConnToRemoteSigner(t *testing.T) {
|
||||
thisConnTimeout = testTimeoutReadWrite
|
||||
validatorEndpoint = newSignerValidatorEndpoint(logger, tc.addr, thisConnTimeout)
|
||||
)
|
||||
// Ping every:
|
||||
SignerValidatorEndpointSetHeartbeat(testTimeoutHeartbeat)(validatorEndpoint)
|
||||
|
||||
SignerServiceEndpointTimeoutReadWrite(testTimeoutReadWrite)(serviceEndpoint)
|
||||
SignerServiceEndpointConnRetries(10)(serviceEndpoint)
|
||||
@@ -110,7 +105,6 @@ func TestRetryConnToRemoteSigner(t *testing.T) {
|
||||
assert.True(t, serviceEndpoint.IsRunning())
|
||||
|
||||
<-readyCh
|
||||
time.Sleep(testTimeoutHeartbeat * 2)
|
||||
|
||||
serviceEndpoint.Stop()
|
||||
rs2 := NewSignerDialerEndpoint(
|
||||
@@ -120,7 +114,6 @@ func TestRetryConnToRemoteSigner(t *testing.T) {
|
||||
tc.dialer,
|
||||
)
|
||||
// let some pings pass
|
||||
time.Sleep(testTimeoutHeartbeat3o2)
|
||||
require.NoError(t, rs2.Start())
|
||||
assert.True(t, rs2.IsRunning())
|
||||
defer rs2.Stop()
|
||||
@@ -193,7 +186,6 @@ func getMockEndpoints(
|
||||
validatorEndpoint = newSignerValidatorEndpoint(logger, addr, testTimeoutReadWrite)
|
||||
)
|
||||
|
||||
SignerValidatorEndpointSetHeartbeat(testTimeoutHeartbeat)(validatorEndpoint)
|
||||
SignerServiceEndpointTimeoutReadWrite(testTimeoutReadWrite)(serviceEndpoint)
|
||||
SignerServiceEndpointConnRetries(1e6)(serviceEndpoint)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user