lotus/node/modules/storageminer.go

160 lines
4.2 KiB
Go
Raw Normal View History

2019-08-01 14:19:53 +00:00
package modules
import (
"context"
2019-08-08 04:24:49 +00:00
"path/filepath"
2019-08-06 22:04:21 +00:00
"github.com/ipfs/go-bitswap"
"github.com/ipfs/go-bitswap/network"
"github.com/ipfs/go-blockservice"
2019-08-01 14:19:53 +00:00
"github.com/ipfs/go-datastore"
2019-08-06 22:04:21 +00:00
blockstore "github.com/ipfs/go-ipfs-blockstore"
"github.com/ipfs/go-merkledag"
2019-08-01 14:19:53 +00:00
"github.com/libp2p/go-libp2p-core/host"
2019-08-29 15:09:34 +00:00
"github.com/libp2p/go-libp2p-core/routing"
2019-08-01 14:19:53 +00:00
"github.com/mitchellh/go-homedir"
"go.uber.org/fx"
"github.com/filecoin-project/lotus/api"
"github.com/filecoin-project/lotus/build"
"github.com/filecoin-project/lotus/chain/address"
"github.com/filecoin-project/lotus/chain/deals"
"github.com/filecoin-project/lotus/lib/sectorbuilder"
"github.com/filecoin-project/lotus/node/modules/dtypes"
"github.com/filecoin-project/lotus/node/modules/helpers"
"github.com/filecoin-project/lotus/node/repo"
"github.com/filecoin-project/lotus/retrieval"
"github.com/filecoin-project/lotus/storage"
"github.com/filecoin-project/lotus/storage/commitment"
"github.com/filecoin-project/lotus/storage/sector"
2019-08-01 14:19:53 +00:00
)
func minerAddrFromDS(ds dtypes.MetadataDS) (address.Address, error) {
maddrb, err := ds.Get(datastore.NewKey("miner-address"))
if err != nil {
return address.Undef, err
}
return address.NewFromBytes(maddrb)
}
func SectorBuilderConfig(storagePath string) func(dtypes.MetadataDS) (*sectorbuilder.SectorBuilderConfig, error) {
return func(ds dtypes.MetadataDS) (*sectorbuilder.SectorBuilderConfig, error) {
minerAddr, err := minerAddrFromDS(ds)
if err != nil {
return nil, err
}
2019-08-01 14:19:53 +00:00
sp, err := homedir.Expand(storagePath)
if err != nil {
return nil, err
}
metadata := filepath.Join(sp, "meta")
sealed := filepath.Join(sp, "sealed")
staging := filepath.Join(sp, "staging")
sb := &sectorbuilder.SectorBuilderConfig{
Miner: minerAddr,
2019-08-27 19:54:14 +00:00
SectorSize: build.SectorSize,
2019-08-01 14:19:53 +00:00
MetadataDir: metadata,
SealedDir: sealed,
StagedDir: staging,
}
return sb, nil
}
}
2019-09-16 16:40:26 +00:00
func StorageMiner(mctx helpers.MetricsCtx, lc fx.Lifecycle, api api.FullNode, h host.Host, ds dtypes.MetadataDS, secst *sector.Store, commt *commitment.Tracker) (*storage.Miner, error) {
maddr, err := minerAddrFromDS(ds)
2019-08-01 14:19:53 +00:00
if err != nil {
return nil, err
}
2019-09-16 16:40:26 +00:00
sm, err := storage.NewMiner(api, maddr, h, ds, secst, commt)
2019-08-01 14:19:53 +00:00
if err != nil {
return nil, err
}
ctx := helpers.LifecycleCtx(mctx, lc)
lc.Append(fx.Hook{
OnStart: func(context.Context) error {
return sm.Run(ctx)
},
})
return sm, nil
}
2019-08-02 16:25:10 +00:00
2019-08-26 13:45:36 +00:00
func HandleRetrieval(host host.Host, lc fx.Lifecycle, m *retrieval.Miner) {
lc.Append(fx.Hook{
OnStart: func(context.Context) error {
2019-08-27 18:45:21 +00:00
host.SetStreamHandler(retrieval.QueryProtocolID, m.HandleQueryStream)
host.SetStreamHandler(retrieval.ProtocolID, m.HandleDealStream)
2019-08-26 13:45:36 +00:00
return nil
},
})
}
2019-08-06 23:08:34 +00:00
func HandleDeals(mctx helpers.MetricsCtx, lc fx.Lifecycle, host host.Host, h *deals.Handler) {
ctx := helpers.LifecycleCtx(mctx, lc)
lc.Append(fx.Hook{
OnStart: func(context.Context) error {
h.Run(ctx)
host.SetStreamHandler(deals.ProtocolID, h.HandleStream)
2019-09-13 21:00:36 +00:00
host.SetStreamHandler(deals.AskProtocolID, h.HandleAskStream)
2019-08-06 23:08:34 +00:00
return nil
},
OnStop: func(context.Context) error {
h.Stop()
return nil
},
})
2019-08-02 16:25:10 +00:00
}
2019-08-06 22:04:21 +00:00
func StagingDAG(mctx helpers.MetricsCtx, lc fx.Lifecycle, r repo.LockedRepo, rt routing.Routing, h host.Host) (dtypes.StagingDAG, error) {
stagingds, err := r.Datastore("/staging")
if err != nil {
return nil, err
}
bs := blockstore.NewBlockstore(stagingds)
ibs := blockstore.NewIdStore(bs)
bitswapNetwork := network.NewFromIpfsHost(h, rt)
exch := bitswap.New(helpers.LifecycleCtx(mctx, lc), bitswapNetwork, bs)
bsvc := blockservice.New(ibs, exch)
dag := merkledag.NewDAGService(bsvc)
lc.Append(fx.Hook{
OnStop: func(_ context.Context) error {
return bsvc.Close()
},
})
return dag, nil
}
func RegisterMiner(lc fx.Lifecycle, ds dtypes.MetadataDS, api api.FullNode) error {
minerAddr, err := minerAddrFromDS(ds)
if err != nil {
return err
}
lc.Append(fx.Hook{
OnStart: func(ctx context.Context) error {
2019-09-17 14:23:08 +00:00
log.Infof("Registering miner '%s' with full node", minerAddr)
return api.MinerRegister(ctx, minerAddr)
},
2019-09-17 14:23:08 +00:00
OnStop: func(ctx context.Context) error {
log.Infof("Unregistering miner '%s' from full node", minerAddr)
return api.MinerUnregister(ctx, minerAddr)
},
})
return nil
}