add routine to tag miners peer IDs as high value in connection manager
This commit is contained in:
@@ -40,6 +40,7 @@ import (
|
||||
"github.com/filecoin-project/lotus/node/modules/helpers"
|
||||
"github.com/filecoin-project/lotus/node/modules/lp2p"
|
||||
"github.com/filecoin-project/lotus/node/modules/testing"
|
||||
"github.com/filecoin-project/lotus/node/peers"
|
||||
"github.com/filecoin-project/lotus/node/repo"
|
||||
"github.com/filecoin-project/lotus/paych"
|
||||
"github.com/filecoin-project/lotus/peermgr"
|
||||
@@ -234,6 +235,7 @@ func Online() Option {
|
||||
Override(new(*paych.Store), paych.NewStore),
|
||||
Override(new(*paych.Manager), paych.NewManager),
|
||||
Override(new(*market.FundMgr), market.NewFundMgr),
|
||||
Override(new(*peers.PeerTagger), lp2p.TagMiners),
|
||||
),
|
||||
|
||||
// Storage miner
|
||||
|
||||
+5
-2
@@ -2,6 +2,7 @@ package impl
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/gbrlsnchs/jwt/v3"
|
||||
"github.com/libp2p/go-libp2p-core/host"
|
||||
"github.com/libp2p/go-libp2p-core/network"
|
||||
@@ -13,13 +14,15 @@ import (
|
||||
"github.com/filecoin-project/lotus/api"
|
||||
"github.com/filecoin-project/lotus/build"
|
||||
"github.com/filecoin-project/lotus/node/modules/dtypes"
|
||||
"github.com/filecoin-project/lotus/node/peers"
|
||||
)
|
||||
|
||||
type CommonAPI struct {
|
||||
fx.In
|
||||
|
||||
APISecret *dtypes.APIAlg
|
||||
Host host.Host
|
||||
APISecret *dtypes.APIAlg
|
||||
Host host.Host
|
||||
PeerTagger *peers.PeerTagger // TODO: this needs a better home
|
||||
}
|
||||
|
||||
type jwtPayload struct {
|
||||
|
||||
+1
-12
@@ -223,18 +223,7 @@ func (a *StateAPI) StateGetReceipt(ctx context.Context, msg cid.Cid, ts *types.T
|
||||
}
|
||||
|
||||
func (a *StateAPI) StateListMiners(ctx context.Context, ts *types.TipSet) ([]address.Address, error) {
|
||||
var state actors.StoragePowerState
|
||||
if _, err := a.StateManager.LoadActorState(ctx, actors.StoragePowerAddress, &state, ts); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
cst := hamt.CSTFromBstore(a.StateManager.ChainStore().Blockstore())
|
||||
miners, err := actors.MinerSetList(ctx, cst, state.Miners)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return miners, nil
|
||||
return stmgr.ListMinerActors(ctx, a.StateManager, ts)
|
||||
}
|
||||
|
||||
func (a *StateAPI) StateListActors(ctx context.Context, ts *types.TipSet) ([]address.Address, error) {
|
||||
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
mocknet "github.com/libp2p/go-libp2p/p2p/net/mock"
|
||||
"go.uber.org/fx"
|
||||
|
||||
"github.com/filecoin-project/lotus/build"
|
||||
"github.com/filecoin-project/lotus/node/modules/dtypes"
|
||||
"github.com/filecoin-project/lotus/node/modules/helpers"
|
||||
)
|
||||
@@ -41,7 +42,7 @@ func Host(mctx helpers.MetricsCtx, lc fx.Lifecycle, params P2PHostIn) (RawHost,
|
||||
return nil, fmt.Errorf("missing private key for node ID: %s", params.ID.Pretty())
|
||||
}
|
||||
|
||||
opts := []libp2p.Option{libp2p.Identity(pkey), libp2p.Peerstore(params.Peerstore), libp2p.NoListenAddrs}
|
||||
opts := []libp2p.Option{libp2p.Identity(pkey), libp2p.Peerstore(params.Peerstore), libp2p.NoListenAddrs, libp2p.UserAgent("lotus-" + build.Version)}
|
||||
for _, o := range params.Opts {
|
||||
opts = append(opts, o...)
|
||||
}
|
||||
|
||||
@@ -1,16 +1,20 @@
|
||||
package lp2p
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"time"
|
||||
|
||||
"github.com/filecoin-project/lotus/chain/stmgr"
|
||||
"github.com/filecoin-project/lotus/chain/types"
|
||||
"github.com/filecoin-project/lotus/node/peers"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
logging "github.com/ipfs/go-log"
|
||||
"github.com/libp2p/go-libp2p"
|
||||
connmgr "github.com/libp2p/go-libp2p-connmgr"
|
||||
"github.com/libp2p/go-libp2p-core/crypto"
|
||||
host "github.com/libp2p/go-libp2p-core/host"
|
||||
"github.com/libp2p/go-libp2p-core/peer"
|
||||
"github.com/libp2p/go-libp2p-core/peerstore"
|
||||
"go.uber.org/fx"
|
||||
@@ -95,3 +99,18 @@ func simpleOpt(opt libp2p.Option) func() (opts Libp2pOpts, err error) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func TagMiners(lc fx.Lifecycle, h host.Host, stmgr *stmgr.StateManager) *peers.PeerTagger {
|
||||
pt := peers.NewPeerTagger(h, stmgr)
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(ctx context.Context) error {
|
||||
pt.Run()
|
||||
return nil
|
||||
},
|
||||
OnStop: func(ctx context.Context) error {
|
||||
return pt.Close()
|
||||
},
|
||||
})
|
||||
|
||||
return pt
|
||||
}
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package peers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/filecoin-project/lotus/chain/stmgr"
|
||||
host "github.com/libp2p/go-libp2p-core/host"
|
||||
"github.com/prometheus/common/log"
|
||||
"go.opencensus.io/trace"
|
||||
)
|
||||
|
||||
// PeerTagger uses information from the chain to tag peer connections to
|
||||
// prevent them from being closed by the connection manager
|
||||
type PeerTagger struct {
|
||||
h host.Host
|
||||
st *stmgr.StateManager
|
||||
|
||||
closing chan struct{}
|
||||
}
|
||||
|
||||
func NewPeerTagger(h host.Host, st *stmgr.StateManager) *PeerTagger {
|
||||
log.Error("NEW PEER TAGGER")
|
||||
return &PeerTagger{
|
||||
h: h,
|
||||
st: st,
|
||||
closing: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
func (pt *PeerTagger) Run() {
|
||||
go pt.loop()
|
||||
}
|
||||
|
||||
func (pt *PeerTagger) loop() {
|
||||
tick := time.NewTimer(time.Second)
|
||||
for {
|
||||
select {
|
||||
case <-tick.C:
|
||||
if err := pt.tagMiners(context.TODO()); err != nil {
|
||||
log.Warn("failed to run tag miners: ", err)
|
||||
}
|
||||
|
||||
tick.Reset(time.Minute * 2)
|
||||
|
||||
case <-pt.closing:
|
||||
log.Warn("peer tagger shutting down")
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (pt *PeerTagger) tagMiners(ctx context.Context) error {
|
||||
_, span := trace.StartSpan(ctx, "tagMiners")
|
||||
defer span.End()
|
||||
|
||||
ts := pt.st.ChainStore().GetHeaviestTipSet()
|
||||
mactors, err := stmgr.ListMinerActors(context.TODO(), pt.st, ts)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, m := range mactors {
|
||||
mpow, _, err := stmgr.GetPower(context.TODO(), pt.st, ts, m)
|
||||
if err != nil {
|
||||
log.Warn("failed to get miners power: ", err)
|
||||
continue
|
||||
}
|
||||
if !mpow.IsZero() {
|
||||
pid, err := stmgr.GetMinerPeerID(context.TODO(), pt.st, ts, m)
|
||||
if err != nil {
|
||||
log.Warn("failed to get peer ID for miner: ", err)
|
||||
continue
|
||||
}
|
||||
pt.h.ConnManager().TagPeer(pid, "miner", 10)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pt *PeerTagger) Close() error {
|
||||
close(pt.closing)
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user