WorkerCount on storageminer config

This commit is contained in:
Łukasz Magiera
2019-11-12 18:59:38 +01:00
parent 7514d3e1bd
commit 73ab6c0c66
13 changed files with 80 additions and 52 deletions
+23 -19
View File
@@ -181,7 +181,7 @@ func Online() Option {
// Full node
ApplyIf(isType(repo.RepoFullNode),
ApplyIf(isType(repo.FullNode),
// TODO: Fix offline mode
Override(new(dtypes.BootstrapPeers), modules.BuiltinBootstrap),
@@ -235,7 +235,7 @@ func Online() Option {
),
// Storage miner
ApplyIf(func(s *Settings) bool { return s.nodeType == repo.RepoStorageMiner },
ApplyIf(func(s *Settings) bool { return s.nodeType == repo.StorageMiner },
Override(new(*sectorbuilder.SectorBuilder), modules.SectorBuilder),
Override(new(*sectorblocks.SectorBlocks), sectorblocks.NewSectorBlocks),
Override(new(storage.TicketFn), modules.SealTicketGen),
@@ -266,7 +266,7 @@ func StorageMiner(out *api.StorageMiner) Option {
),
func(s *Settings) error {
s.nodeType = repo.RepoStorageMiner
s.nodeType = repo.StorageMiner
return nil
},
@@ -294,30 +294,34 @@ func ConfigCommon(cfg *config.Common) Option {
)
}
func ConfigFullNode(cfg *config.FullNode) Option {
//ApplyIf(func(s *Settings) bool { return s.nodeType == repo.RepoFullNode }),
func ConfigFullNode(c interface{}) Option {
cfg, ok := c.(*config.FullNode)
if !ok {
return Error(xerrors.Errorf("invalid config from repo, got: %T", c))
}
return Options(
ConfigCommon(&cfg.Common),
Override(HeadMetricsKey, metrics.SendHeadNotifs(cfg.Metrics.Nickname)),
)
}
func configFull(c interface{}) Option {
cfg, ok := c.(*config.FullNode)
if !ok {
return Error(xerrors.Errorf("invalid config from repo, got: %T", c))
}
return ConfigFullNode(cfg)
}
func configMiner(c interface{}) Option {
func ConfigStorageMiner(c interface{}, lr repo.LockedRepo) Option {
cfg, ok := c.(*config.StorageMiner)
if !ok {
return Error(xerrors.Errorf("invalid config from repo, got: %T", c))
}
return ConfigCommon(&cfg.Common)
path := cfg.SectorBuilder.Path
if path == "" {
path = lr.Path()
}
return Options(
ConfigCommon(&cfg.Common),
Override(new(*sectorbuilder.Config), modules.SectorBuilderConfig(path, cfg.SectorBuilder.WorkerCount)),
)
}
func Repo(r repo.Repo) Option {
@@ -334,8 +338,8 @@ func Repo(r repo.Repo) Option {
return Options(
Override(new(repo.LockedRepo), modules.LockedRepo(lr)), // module handles closing
ApplyIf(isType(repo.RepoFullNode), configFull(c)),
ApplyIf(isType(repo.RepoStorageMiner), configMiner(c)),
ApplyIf(isType(repo.FullNode), ConfigFullNode(c)),
ApplyIf(isType(repo.StorageMiner), ConfigStorageMiner(c, lr)),
Override(new(dtypes.MetadataDS), modules.Datastore),
Override(new(dtypes.ChainBlockstore), modules.ChainBlockstore),
@@ -371,7 +375,7 @@ func New(ctx context.Context, opts ...Option) (StopFunc, error) {
settings := Settings{
modules: map[interface{}]fx.Option{},
invokes: make([]fx.Option, _nInvokes),
nodeType: repo.RepoFullNode,
nodeType: repo.FullNode,
}
// apply module options in the right order
+17
View File
@@ -17,9 +17,13 @@ type FullNode struct {
Metrics Metrics
}
// // Common
// StorageMiner is a storage miner config
type StorageMiner struct {
Common
SectorBuilder SectorBuilder
}
// API contains configs for API endpoint
@@ -34,10 +38,19 @@ type Libp2p struct {
BootstrapPeers []string
}
// // Full Node
type Metrics struct {
Nickname string
}
// // Storage Miner
type SectorBuilder struct {
Path string
WorkerCount uint
}
func defCommon() Common {
return Common{
API: API{
@@ -64,6 +77,10 @@ func DefaultFullNode() *FullNode {
func DefaultStorageMiner() *StorageMiner {
return &StorageMiner{
Common: defCommon(),
SectorBuilder: SectorBuilder{
WorkerCount: 5,
},
}
}
+1 -9
View File
@@ -4,9 +4,7 @@ import (
"bytes"
"context"
"crypto/rand"
"io/ioutil"
"net/http/httptest"
"os"
"testing"
"github.com/libp2p/go-libp2p-core/crypto"
@@ -23,7 +21,6 @@ import (
"github.com/filecoin-project/lotus/chain/address"
"github.com/filecoin-project/lotus/chain/types"
"github.com/filecoin-project/lotus/lib/jsonrpc"
"github.com/filecoin-project/lotus/lib/sectorbuilder"
"github.com/filecoin-project/lotus/miner"
"github.com/filecoin-project/lotus/node"
"github.com/filecoin-project/lotus/node/modules"
@@ -34,7 +31,7 @@ import (
func testStorageNode(ctx context.Context, t *testing.T, waddr address.Address, act address.Address, pk crypto.PrivKey, tnd test.TestNode, mn mocknet.Mocknet) test.TestStorageNode {
r := repo.NewMemory(nil)
lr, err := r.Lock(repo.RepoStorageMiner)
lr, err := r.Lock(repo.StorageMiner)
require.NoError(t, err)
ks, err := lr.KeyStore()
@@ -77,10 +74,6 @@ func testStorageNode(ctx context.Context, t *testing.T, waddr address.Address, a
require.NoError(t, err)
// start node
secbpath, err := ioutil.TempDir(os.TempDir(), "lotust-stortest-sb-")
require.NoError(t, err)
var minerapi api.StorageMiner
// TODO: use stop
@@ -92,7 +85,6 @@ func testStorageNode(ctx context.Context, t *testing.T, waddr address.Address, a
node.MockHost(mn),
node.Override(new(*sectorbuilder.Config), modules.SectorBuilderConfig(secbpath, 2)),
node.Override(new(api.FullNode), tnd),
)
require.NoError(t, err)
+5 -5
View File
@@ -37,16 +37,16 @@ const (
type RepoType int
const (
_ = iota // Default is invalid
RepoFullNode RepoType = iota
RepoStorageMiner
_ = iota // Default is invalid
FullNode RepoType = iota
StorageMiner
)
func defConfForType(t RepoType) interface{} {
switch t {
case RepoFullNode:
case FullNode:
return config.DefaultFullNode()
case RepoStorageMiner:
case StorageMiner:
return config.DefaultStorageMiner()
default:
panic(fmt.Sprintf("unknown RepoType(%d)", int(t)))
+1 -1
View File
@@ -17,7 +17,7 @@ func genFsRepo(t *testing.T) (*FsRepo, func()) {
t.Fatal(err)
}
err = repo.Init(RepoFullNode)
err = repo.Init(FullNode)
if err != ErrRepoExists && err != nil {
t.Fatal(err)
}
+20 -2
View File
@@ -1,6 +1,8 @@
package repo
import (
"io/ioutil"
"os"
"sync"
"github.com/ipfs/go-datastore"
@@ -32,11 +34,20 @@ type lockedMemRepo struct {
t RepoType
sync.RWMutex
token *byte
tempDir string
token *byte
}
func (lmem *lockedMemRepo) Path() string {
return ""
t, err := ioutil.TempDir(os.TempDir(), "lotus-memrepo-temp-")
if err != nil {
panic(err) // only used in tests, probably fine
}
lmem.Lock()
lmem.tempDir = t
lmem.Unlock()
return t
}
var _ Repo = &MemRepo{}
@@ -127,6 +138,13 @@ func (lmem *lockedMemRepo) Close() error {
return ErrClosedRepo
}
if lmem.tempDir != "" {
if err := os.RemoveAll(lmem.tempDir); err != nil {
return err
}
lmem.tempDir = ""
}
lmem.mem.token = nil
lmem.mem.api.Lock()
lmem.mem.api.ma = nil
+4 -4
View File
@@ -18,12 +18,12 @@ func basicTest(t *testing.T, repo Repo) {
}
assert.Nil(t, apima, "with no api endpoint, return should be nil")
lrepo, err := repo.Lock(RepoFullNode)
lrepo, err := repo.Lock(FullNode)
assert.NoError(t, err, "should be able to lock once")
assert.NotNil(t, lrepo, "locked repo shouldn't be nil")
{
lrepo2, err := repo.Lock(RepoFullNode)
lrepo2, err := repo.Lock(FullNode)
if assert.Error(t, err) {
assert.Equal(t, ErrRepoAlreadyLocked, err)
}
@@ -33,7 +33,7 @@ func basicTest(t *testing.T, repo Repo) {
err = lrepo.Close()
assert.NoError(t, err, "should be able to unlock")
lrepo, err = repo.Lock(RepoFullNode)
lrepo, err = repo.Lock(FullNode)
assert.NoError(t, err, "should be able to relock")
assert.NotNil(t, lrepo, "locked repo shouldn't be nil")
@@ -64,7 +64,7 @@ func basicTest(t *testing.T, repo Repo) {
k1 := types.KeyInfo{Type: "foo"}
k2 := types.KeyInfo{Type: "bar"}
lrepo, err = repo.Lock(RepoFullNode)
lrepo, err = repo.Lock(FullNode)
assert.NoError(t, err, "should be able to relock")
assert.NotNil(t, lrepo, "locked repo shouldn't be nil")