Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion cmd/custom/node.go
Original file line number Diff line number Diff line change
Expand Up @@ -209,7 +209,7 @@ func main() {
zap.S().Info("Successfully dropped peers storage")
}

peerManager := peer_manager.NewPeerManager(peerSpawnerImpl, peerStorage, int(limitConnections), version)
peerManager := peer_manager.NewPeerManager(peerSpawnerImpl, peerStorage, int(limitConnections), version, conf.WavesNetwork)
go peerManager.Run(ctx)

scheduler := scheduler2.NewScheduler(
Expand Down
1 change: 1 addition & 0 deletions cmd/node/node.go
Original file line number Diff line number Diff line change
Expand Up @@ -319,6 +319,7 @@ func main() {
peerStorage,
int(limitConnections),
version,
conf.WavesNetwork,
)
go peerManager.Run(ctx)

Expand Down
298 changes: 140 additions & 158 deletions pkg/mock/peer_manager.go

Large diffs are not rendered by default.

1,859 changes: 929 additions & 930 deletions pkg/mock/state.go

Large diffs are not rendered by default.

15 changes: 7 additions & 8 deletions pkg/node/node.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,12 +101,12 @@ func (a *Node) Serve(ctx context.Context) error {
}
}

func (a *Node) logErrors(err error) {
func (a *Node) logErrors(fsm state_fsm.FSM, err error) {
switch e := err.(type) {
case *proto.InfoMsg:
zap.S().Debug(e.Error())
zap.S().Debugf("[%s] %s", fsm.String(), e.Error())
default:
zap.S().Error(e.Error())
zap.S().Errorf("[%s] %s", fsm.String(), e.Error())
}
}

Expand Down Expand Up @@ -165,7 +165,7 @@ func (a *Node) Run(ctx context.Context, p peer.Parent, InternalMessageCh chan me
default:
}
default:
zap.S().Errorf("unknown internalMess %T", t)
zap.S().Errorf("[%s] Unknown internal message '%T'", fsm.String(), t)
continue
}
case task := <-tasksCh:
Expand All @@ -178,19 +178,18 @@ func (a *Node) Run(ctx context.Context, p peer.Parent, InternalMessageCh chan me
fsm, async, err = fsm.PeerError(m.Peer, t)
}
case mess := <-p.MessageCh:
zap.S().Debugf("received proto Message %T", mess.Message)
zap.S().Debugf("[%s] Network message '%T' received", fsm.String(), mess.Message)
action, ok := actions[reflect.TypeOf(mess.Message)]
if !ok {
zap.S().Errorf("unknown proto Message %T", mess.Message)
zap.S().Errorf("[%s] Unknown network message '%T'", fsm.String(), mess.Message)
continue
}
fsm, async, err = action(a.services, mess, fsm)
}
if err != nil {
a.logErrors(err)
a.logErrors(fsm, err)
}
spawnAsync(ctx, tasksCh, a.services.LoggableRunner, async)
zap.S().Debugf("FSM %T", fsm)
}
}

Expand Down
47 changes: 15 additions & 32 deletions pkg/node/peer_manager/peer_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,13 @@ package peer_manager

import (
"context"
"github.com/wavesplatform/gowaves/pkg/node/peer_manager/storage"
"math/big"
"net"
"sort"
"sync"
"time"

"github.com/wavesplatform/gowaves/pkg/node/peer_manager/storage"

"github.com/pkg/errors"
"github.com/wavesplatform/gowaves/pkg/p2p/peer"
"github.com/wavesplatform/gowaves/pkg/proto"
Expand All @@ -29,12 +29,6 @@ func newPeerInfo(peer peer.Peer) peerInfo {
}
}

type byScore []peerInfo

func (a byScore) Len() int { return len(a) }
func (a byScore) Less(i, j int) bool { return a[i].score.Cmp(a[j].score) < 0 }
func (a byScore) Swap(i, j int) { a[i], a[j] = a[j], a[i] }

type PeerManager interface {
Connected(peer.Peer) (peer.Peer, bool)
NewConnection(peer.Peer) error
Expand All @@ -45,7 +39,6 @@ type PeerManager interface {
Suspend(peer peer.Peer, suspendTime time.Time, reason string)
Suspended() []storage.SuspendedPeer
AddConnected(peer.Peer)
PeerWithHighestScore() (peer.Peer, *big.Int, bool)
UpdateScore(p peer.Peer, score *proto.Score) error
UpdateKnownPeers([]storage.KnownPeer) error
KnownPeers() []storage.KnownPeer
Expand All @@ -71,10 +64,11 @@ type PeerManagerImpl struct {
connectPeers bool // spawn outgoing
limitConnections int
version proto.Version
networkName string
}

func NewPeerManager(spawner PeerSpawner, storage PeerStorage,
limitConnections int, version proto.Version) *PeerManagerImpl {
limitConnections int, version proto.Version, networkName string) *PeerManagerImpl {

return &PeerManagerImpl{
spawner: spawner,
Expand All @@ -84,6 +78,7 @@ func NewPeerManager(spawner PeerSpawner, storage PeerStorage,
connectPeers: true,
limitConnections: limitConnections,
version: version,
networkName: networkName,
}
}

Expand Down Expand Up @@ -130,7 +125,7 @@ func (a *PeerManagerImpl) NewConnection(p peer.Peer) error {
}
if a.IsSuspended(p) {
_ = p.Close()
return errors.New("peer is suspended")
return errors.Errorf("peer '%s' is suspended", p.ID())
}
if p.Handshake().Version.CmpMinor(a.version) >= 2 {
err := errors.Errorf(
Expand All @@ -142,7 +137,13 @@ func (a *PeerManagerImpl) NewConnection(p peer.Peer) error {
_ = p.Close()
return proto.NewInfoMsg(err)
}

if p.Handshake().AppName != a.networkName {
err := errors.Errorf("peer '%s' has the invalid network name '%s', required '%s'",
p.ID(), p.Handshake().AppName, a.networkName)
a.Suspend(p, time.Now(), err.Error())
_ = p.Close()
return proto.NewInfoMsg(err)
}
in, out := a.InOutCount()
switch p.Direction() {
case peer.Incoming:
Expand Down Expand Up @@ -194,33 +195,15 @@ func (a *PeerManagerImpl) AddConnected(peer peer.Peer) {
a.active[peer] = newPeerInfo(peer)
}

func (a *PeerManagerImpl) PeerWithHighestScore() (peer.Peer, *big.Int, bool) {
a.mu.RLock()
defer a.mu.RUnlock()

if len(a.active) == 0 {
return nil, nil, false
}

peers := make([]peerInfo, 0)
for _, p := range a.active {
peers = append(peers, p)
}

sort.Sort(byScore(peers))

highest := peers[len(peers)-1]
return highest.peer, highest.score, true
}

func (a *PeerManagerImpl) UpdateScore(p peer.Peer, score *big.Int) error {
a.mu.Lock()
defer a.mu.Unlock()
if row, ok := a.active[p]; ok {
row.score = score
a.active[p] = row
return nil
}
return nil
return errors.Errorf("peer '%s' is not active", p.ID())
}

func (a *PeerManagerImpl) IsSuspended(p peer.Peer) bool {
Expand Down
2 changes: 2 additions & 0 deletions pkg/node/state_fsm/fsm.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,8 @@ type FSM interface {
Transaction(p peer.Peer, t proto.Transaction) (FSM, Async, error)

Halt() (FSM, Async, error)

String() string
}

func NewFsm(
Expand Down
35 changes: 4 additions & 31 deletions pkg/node/state_fsm/fsm_common.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,13 @@ package state_fsm

import (
"github.com/wavesplatform/gowaves/pkg/node/peer_manager"
. "github.com/wavesplatform/gowaves/pkg/p2p/peer"
"github.com/wavesplatform/gowaves/pkg/p2p/peer"
"github.com/wavesplatform/gowaves/pkg/proto"
"github.com/wavesplatform/gowaves/pkg/state"
"go.uber.org/zap"
)

func newPeer(fsm FSM, p Peer, peers peer_manager.PeerManager) (FSM, Async, error) {
func newPeer(fsm FSM, p peer.Peer, peers peer_manager.PeerManager) (FSM, Async, error) {
err := peers.NewConnection(p)
if err != nil {
return fsm, nil, proto.NewInfoMsg(err)
Expand All @@ -17,7 +17,7 @@ func newPeer(fsm FSM, p Peer, peers peer_manager.PeerManager) (FSM, Async, error
}

// TODO handle no peers
func peerError(fsm FSM, p Peer, peers peer_manager.PeerManager, _ error) (FSM, Async, error) {
func peerError(fsm FSM, p peer.Peer, peers peer_manager.PeerManager, _ error) (FSM, Async, error) {
peers.Disconnect(p)
return fsm, nil, nil
}
Expand All @@ -26,24 +26,7 @@ func noop(fsm FSM) (FSM, Async, error) {
return fsm, nil, nil
}

func handleScore(fsm FSM, info BaseInfo, p Peer, score *proto.Score) (FSM, Async, error) {
err := info.peers.UpdateScore(p, score)
if err != nil {
return fsm, nil, err
}

myScore, err := info.storage.CurrentScore()
if err != nil {
return NewIdleFsm(info), nil, err
}

if score.Cmp(myScore) == 1 { // remote score > my score
return NewIdleToSyncTransition(info, p)
}
return fsm, nil, nil
}

func sendScore(p Peer, storage state.State) {
func sendScore(p peer.Peer, storage state.State) {
curScore, err := storage.CurrentScore()
if err != nil {
zap.S().Error(err)
Expand All @@ -53,13 +36,3 @@ func sendScore(p Peer, storage state.State) {
bts := curScore.Bytes()
p.SendMessage(&proto.ScoreMessage{Score: bts})
}

// TODO send micro block
//func handleMineMicro(a FromBaseInfo, base BaseInfo, minedBlock *proto.Block, rest miner.MiningLimits, blocks ng.Blocks, keyPair proto.KeyPair) (FSM, Async, error) {
// block, micro, rest, err := base.microMiner.Micro(rest, minedBlock, blocks, keyPair)
// if err != nil {
// return a, nil, err
// }
// base.
// return a.FromBaseInfo()
//}
4 changes: 4 additions & 0 deletions pkg/node/state_fsm/fsm_halt.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,10 @@ func (a HaltFSM) MicroBlockInv(p peer.Peer, inv *proto.MicroBlockInv) (FSM, Asyn
return noop(a)
}

func (a HaltFSM) String() string {
return "Halt"
}

func HaltTransition(info BaseInfo) (FSM, Async, error) {
zap.S().Debugf("started HaltTransition ")
info.peers.Close()
Expand Down
57 changes: 43 additions & 14 deletions pkg/node/state_fsm/fsm_idle.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,15 @@
package state_fsm

import (
"time"

"github.com/pkg/errors"
"github.com/wavesplatform/gowaves/pkg/libs/signatures"
"github.com/wavesplatform/gowaves/pkg/metrics"
. "github.com/wavesplatform/gowaves/pkg/node/state_fsm/tasks"
"github.com/wavesplatform/gowaves/pkg/node/state_fsm/sync_internal"
"github.com/wavesplatform/gowaves/pkg/node/state_fsm/tasks"
"github.com/wavesplatform/gowaves/pkg/p2p/peer"
"github.com/wavesplatform/gowaves/pkg/p2p/peer/extension"
"github.com/wavesplatform/gowaves/pkg/proto"
"github.com/wavesplatform/gowaves/pkg/types"
"go.uber.org/zap"
Expand Down Expand Up @@ -43,18 +48,18 @@ func (a *IdleFsm) MicroBlockInv(_ peer.Peer, _ *proto.MicroBlockInv) (FSM, Async
return a.baseInfo.d.Noop(a)
}

func (a *IdleFsm) Task(task AsyncTask) (FSM, Async, error) {
zap.S().Debugf("IdleFsm Task: got task type %d, data %+v", task.TaskType, task.Data)
func (a *IdleFsm) Task(task tasks.AsyncTask) (FSM, Async, error) {
switch task.TaskType {
case Ping:
case tasks.Ping:
return noop(a)
case AskPeers:
case tasks.AskPeers:
zap.S().Debug("[Idle] Requesting peers")
a.baseInfo.peers.AskPeers()
return a, nil, nil
case MineMicro: // Do nothing
case tasks.MineMicro: // Do nothing
return a, nil, nil
default:
return a, nil, errors.Errorf("IdleFsm Task: unknown task type %d, data %+v", task.TaskType, task.Data)
return a, nil, errors.Errorf("unexpected internal task '%d' with data '%+v' received by %s FSM", task.TaskType, task.Data, a.String())
}
}

Expand All @@ -66,12 +71,6 @@ func (a *IdleFsm) BlockIDs(_ peer.Peer, _ []proto.BlockID) (FSM, Async, error) {
return a.baseInfo.d.Noop(a)
}

func NewIdleFsm(info BaseInfo) *IdleFsm {
return &IdleFsm{
baseInfo: info,
}
}

func (a *IdleFsm) NewPeer(p peer.Peer) (FSM, Async, error) {
fsm, as, err := newPeer(a, p, a.baseInfo.peers)
if a.baseInfo.peers.ConnectedCount() == a.baseInfo.minPeersMining {
Expand All @@ -83,9 +82,39 @@ func (a *IdleFsm) NewPeer(p peer.Peer) (FSM, Async, error) {

func (a *IdleFsm) Score(p peer.Peer, score *proto.Score) (FSM, Async, error) {
metrics.FSMScore("idle", score, p.Handshake().NodeName)
return handleScore(a, a.baseInfo, p, score)
if err := a.baseInfo.peers.UpdateScore(p, score); err != nil {
return a, nil, err
}
nodeScore, err := a.baseInfo.storage.CurrentScore()
if err != nil {
return a, nil, err
}
if score.Cmp(nodeScore) == 1 {
lastSignatures, err := signatures.LastSignaturesImpl{}.LastBlockIDs(a.baseInfo.storage)
if err != nil {
return a, nil, err
}
internal := sync_internal.InternalFromLastSignatures(extension.NewPeerExtension(p, a.baseInfo.scheme), lastSignatures)
c := conf{
peerSyncWith: p,
timeout: 30 * time.Second,
}
zap.S().Debugf("[Idle] Starting synchronisation with peer '%s'", p.ID())
return NewSyncFsm(a.baseInfo, c.Now(a.baseInfo.tm), internal)
}
return noop(a)
}

func (a *IdleFsm) Block(_ peer.Peer, _ *proto.Block) (FSM, Async, error) {
return noop(a)
}

func (a *IdleFsm) String() string {
return "Idle"
}

func NewIdleFsm(info BaseInfo) *IdleFsm {
return &IdleFsm{
baseInfo: info,
}
}
Loading