From 22ebd6f76fc8e92ceaa6d26f9c6c8b880994f330 Mon Sep 17 00:00:00 2001 From: Dirk McCormick Date: Tue, 10 Aug 2021 15:01:04 +0200 Subject: [PATCH] refactor: extract SectorAccessor from RetrievalProviderNode --- go.mod | 2 +- go.sum | 4 +- markets/dagstore/miner_api.go | 16 +-- markets/dagstore/wrapper.go | 6 +- markets/retrievaladapter/provider.go | 119 ++------------------ markets/sectoraccessor/sectoraccessor.go | 132 +++++++++++++++++++++++ markets/storageadapter/provider.go | 12 +-- node/builder_miner.go | 4 + node/impl/client/client.go | 25 ++++- node/impl/storminer.go | 26 ++--- node/modules/client.go | 8 +- node/modules/storageminer.go | 3 + node/modules/storageminer_dagstore.go | 8 +- 13 files changed, 206 insertions(+), 159 deletions(-) create mode 100644 markets/sectoraccessor/sectoraccessor.go diff --git a/go.mod b/go.mod index 305dd2425..078e6748a 100644 --- a/go.mod +++ b/go.mod @@ -36,7 +36,7 @@ require ( github.com/filecoin-project/go-data-transfer v1.7.3 github.com/filecoin-project/go-fil-commcid v0.1.0 github.com/filecoin-project/go-fil-commp-hashhash v0.1.0 - github.com/filecoin-project/go-fil-markets v1.6.3-0.20210806151421-967d8717c917 + github.com/filecoin-project/go-fil-markets v1.6.3-0.20210810124704-13ab16627deb github.com/filecoin-project/go-jsonrpc v0.1.4-0.20210217175800-45ea43ac2bec github.com/filecoin-project/go-multistore v0.0.3 github.com/filecoin-project/go-padreader v0.0.0-20210723183308-812a16dc01b1 diff --git a/go.sum b/go.sum index 8715e9b2a..063464736 100644 --- a/go.sum +++ b/go.sum @@ -289,8 +289,8 @@ github.com/filecoin-project/go-fil-commcid v0.1.0/go.mod h1:Eaox7Hvus1JgPrL5+M3+ github.com/filecoin-project/go-fil-commp-hashhash v0.1.0 h1:imrrpZWEHRnNqqv0tN7LXep5bFEVOVmQWHJvl2mgsGo= github.com/filecoin-project/go-fil-commp-hashhash v0.1.0/go.mod h1:73S8WSEWh9vr0fDJVnKADhfIv/d6dCbAGaAGWbdJEI8= github.com/filecoin-project/go-fil-markets v1.0.5-0.20201113164554-c5eba40d5335/go.mod h1:AJySOJC00JRWEZzRG2KsfUnqEf5ITXxeX09BE9N4f9c= -github.com/filecoin-project/go-fil-markets v1.6.3-0.20210806151421-967d8717c917 h1:A7B4EShovtH4B8tMp6+Smst7B+H+CLSJ4tcjcS15ta4= -github.com/filecoin-project/go-fil-markets v1.6.3-0.20210806151421-967d8717c917/go.mod h1:9DxT29OFLUJPswugvIbqD5Ab768pJQxkxbnTEFTI8lk= +github.com/filecoin-project/go-fil-markets v1.6.3-0.20210810124704-13ab16627deb h1:+5ODoU3LjQ2Gxi+9Fe9Oe/fYVNG5l5BkHHLpxkrH8LE= +github.com/filecoin-project/go-fil-markets v1.6.3-0.20210810124704-13ab16627deb/go.mod h1:uTcTs8PQv1j72xSrnH3MrXnq3bUE2kQu89CNmRhsRMI= github.com/filecoin-project/go-hamt-ipld v0.1.5 h1:uoXrKbCQZ49OHpsTCkrThPNelC4W3LPEk0OrS/ytIBM= github.com/filecoin-project/go-hamt-ipld v0.1.5/go.mod h1:6Is+ONR5Cd5R6XZoCse1CWaXZc0Hdb/JeX+EQCQzX24= github.com/filecoin-project/go-hamt-ipld/v2 v2.0.0 h1:b3UDemBYN2HNfk3KOXNuxgTTxlWi3xVvbQP0IT38fvM= diff --git a/markets/dagstore/miner_api.go b/markets/dagstore/miner_api.go index 0da01c396..ca6483723 100644 --- a/markets/dagstore/miner_api.go +++ b/markets/dagstore/miner_api.go @@ -10,7 +10,7 @@ import ( "golang.org/x/xerrors" "github.com/filecoin-project/go-fil-markets/piecestore" - "github.com/filecoin-project/go-fil-markets/retrievalmarket" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" "github.com/filecoin-project/go-fil-markets/shared" ) @@ -23,17 +23,17 @@ type MinerAPI interface { type minerAPI struct { pieceStore piecestore.PieceStore - rm retrievalmarket.RetrievalProviderNode + sa sectoraccessor.SectorAccessor throttle throttle.Throttler readyMgr *shared.ReadyManager } var _ MinerAPI = (*minerAPI)(nil) -func NewMinerAPI(store piecestore.PieceStore, rm retrievalmarket.RetrievalProviderNode, concurrency int) MinerAPI { +func NewMinerAPI(store piecestore.PieceStore, sa sectoraccessor.SectorAccessor, concurrency int) MinerAPI { return &minerAPI{ pieceStore: store, - rm: rm, + sa: sa, throttle: throttle.Fixed(concurrency), readyMgr: shared.NewReadyManager(), } @@ -70,7 +70,7 @@ func (m *minerAPI) IsUnsealed(ctx context.Context, pieceCid cid.Cid) (bool, erro var isUnsealed bool // Throttle this path to avoid flooding the storage subsystem. err := m.throttle.Do(ctx, func(ctx context.Context) (err error) { - isUnsealed, err = m.rm.IsUnsealed(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) + isUnsealed, err = m.sa.IsUnsealed(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) if err != nil { return fmt.Errorf("failed to check if sector %d for deal %d was unsealed: %w", deal.SectorID, deal.DealID, err) } @@ -119,7 +119,7 @@ func (m *minerAPI) FetchUnsealedPiece(ctx context.Context, pieceCid cid.Cid) (io // Throttle this path to avoid flooding the storage subsystem. var reader io.ReadCloser err := m.throttle.Do(ctx, func(ctx context.Context) (err error) { - isUnsealed, err := m.rm.IsUnsealed(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) + isUnsealed, err := m.sa.IsUnsealed(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) if err != nil { return fmt.Errorf("failed to check if sector %d for deal %d was unsealed: %w", deal.SectorID, deal.DealID, err) } @@ -127,7 +127,7 @@ func (m *minerAPI) FetchUnsealedPiece(ctx context.Context, pieceCid cid.Cid) (io return nil } // Because we know we have an unsealed copy, this UnsealSector call will actually not perform any unsealing. - reader, err = m.rm.UnsealSector(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) + reader, err = m.sa.UnsealSector(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) return err }) @@ -149,7 +149,7 @@ func (m *minerAPI) FetchUnsealedPiece(ctx context.Context, pieceCid cid.Cid) (io // block for a long time with the current PoRep // // This path is unthrottled. - reader, err := m.rm.UnsealSector(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) + reader, err := m.sa.UnsealSector(ctx, deal.SectorID, deal.Offset.Unpadded(), deal.Length.Unpadded()) if err != nil { lastErr = xerrors.Errorf("failed to unseal deal %d: %w", deal.DealID, err) log.Warn(lastErr.Error()) diff --git a/markets/dagstore/wrapper.go b/markets/dagstore/wrapper.go index 75350c6f4..9d975262f 100644 --- a/markets/dagstore/wrapper.go +++ b/markets/dagstore/wrapper.go @@ -53,10 +53,10 @@ type Wrapper struct { var _ stores.DAGStoreWrapper = (*Wrapper)(nil) -func NewDAGStore(cfg config.DAGStoreConfig, mountApi MinerAPI) (*dagstore.DAGStore, *Wrapper, error) { +func NewDAGStore(cfg config.DAGStoreConfig, minerApi MinerAPI) (*dagstore.DAGStore, *Wrapper, error) { // construct the DAG Store. registry := mount.NewRegistry() - if err := registry.Register(lotusScheme, mountTemplate(mountApi)); err != nil { + if err := registry.Register(lotusScheme, mountTemplate(minerApi)); err != nil { return nil, nil, xerrors.Errorf("failed to create registry: %w", err) } @@ -104,7 +104,7 @@ func NewDAGStore(cfg config.DAGStoreConfig, mountApi MinerAPI) (*dagstore.DAGSto w := &Wrapper{ cfg: cfg, dagst: dagst, - minerAPI: mountApi, + minerAPI: minerApi, failureCh: failureCh, traceCh: traceCh, gcInterval: time.Duration(cfg.GCInterval), diff --git a/markets/retrievaladapter/provider.go b/markets/retrievaladapter/provider.go index 2f6305805..470c1cfc7 100644 --- a/markets/retrievaladapter/provider.go +++ b/markets/retrievaladapter/provider.go @@ -2,44 +2,34 @@ package retrievaladapter import ( "context" - "io" - "github.com/filecoin-project/lotus/api" - "github.com/filecoin-project/lotus/api/v1api" - "github.com/filecoin-project/lotus/node/modules/dtypes" - "github.com/filecoin-project/lotus/storage/sectorblocks" "github.com/hashicorp/go-multierror" "golang.org/x/xerrors" "github.com/ipfs/go-cid" - "github.com/filecoin-project/lotus/chain/actors/builtin/paych" - "github.com/filecoin-project/lotus/chain/types" - sectorstorage "github.com/filecoin-project/lotus/extern/sector-storage" - "github.com/filecoin-project/lotus/extern/sector-storage/storiface" - "github.com/filecoin-project/go-address" "github.com/filecoin-project/go-fil-markets/retrievalmarket" "github.com/filecoin-project/go-fil-markets/shared" "github.com/filecoin-project/go-state-types/abi" - specstorage "github.com/filecoin-project/specs-storage/storage" - + "github.com/filecoin-project/lotus/api/v1api" + "github.com/filecoin-project/lotus/chain/actors/builtin/paych" + "github.com/filecoin-project/lotus/chain/types" logging "github.com/ipfs/go-log/v2" ) var log = logging.Logger("retrievaladapter") type retrievalProviderNode struct { - maddr address.Address - secb sectorblocks.SectorBuilder - pp sectorstorage.PieceProvider - full v1api.FullNode + full v1api.FullNode } +var _ retrievalmarket.RetrievalProviderNode = (*retrievalProviderNode)(nil) + // NewRetrievalProviderNode returns a new node adapter for a retrieval provider that talks to the // Lotus Node -func NewRetrievalProviderNode(maddr dtypes.MinerAddress, secb sectorblocks.SectorBuilder, pp sectorstorage.PieceProvider, full v1api.FullNode) retrievalmarket.RetrievalProviderNode { - return &retrievalProviderNode{address.Address(maddr), secb, pp, full} +func NewRetrievalProviderNode(full v1api.FullNode) retrievalmarket.RetrievalProviderNode { + return &retrievalProviderNode{full: full} } func (rpn *retrievalProviderNode) GetMinerWorkerAddress(ctx context.Context, miner address.Address, tok shared.TipSetToken) (address.Address, error) { @@ -52,42 +42,6 @@ func (rpn *retrievalProviderNode) GetMinerWorkerAddress(ctx context.Context, min return mi.Worker, err } -func (rpn *retrievalProviderNode) UnsealSector(ctx context.Context, sectorID abi.SectorNumber, offset abi.UnpaddedPieceSize, length abi.UnpaddedPieceSize) (io.ReadCloser, error) { - log.Debugf("get sector %d, offset %d, length %d", sectorID, offset, length) - si, err := rpn.sectorsStatus(ctx, sectorID, false) - if err != nil { - return nil, err - } - - mid, err := address.IDFromAddress(rpn.maddr) - if err != nil { - return nil, err - } - - ref := specstorage.SectorRef{ - ID: abi.SectorID{ - Miner: abi.ActorID(mid), - Number: sectorID, - }, - ProofType: si.SealProof, - } - - var commD cid.Cid - if si.CommD != nil { - commD = *si.CommD - } - - // Get a reader for the piece, unsealing the piece if necessary - log.Debugf("read piece in sector %d, offset %d, length %d from miner %d", sectorID, offset, length, mid) - r, unsealed, err := rpn.pp.ReadPiece(ctx, ref, storiface.UnpaddedByteIndex(offset), length, si.Ticket.Value, commD) - if err != nil { - return nil, xerrors.Errorf("failed to unseal piece from sector %d: %w", sectorID, err) - } - _ = unsealed // todo: use - - return r, nil -} - func (rpn *retrievalProviderNode) SavePaymentVoucher(ctx context.Context, paymentChannel address.Address, voucher *paych.SignedVoucher, proof []byte, expectedAmount abi.TokenAmount, tok shared.TipSetToken) (abi.TokenAmount, error) { // TODO: respect the provided TipSetToken (a serialized TipSetKey) when // querying the chain @@ -104,29 +58,6 @@ func (rpn *retrievalProviderNode) GetChainHead(ctx context.Context) (shared.TipS return head.Key().Bytes(), head.Height(), nil } -func (rpn *retrievalProviderNode) IsUnsealed(ctx context.Context, sectorID abi.SectorNumber, offset abi.UnpaddedPieceSize, length abi.UnpaddedPieceSize) (bool, error) { - si, err := rpn.sectorsStatus(ctx, sectorID, true) - if err != nil { - return false, xerrors.Errorf("failed to get sector info: %w", err) - } - - mid, err := address.IDFromAddress(rpn.maddr) - if err != nil { - return false, err - } - - ref := specstorage.SectorRef{ - ID: abi.SectorID{ - Miner: abi.ActorID(mid), - Number: sectorID, - }, - ProofType: si.SealProof, - } - - log.Debugf("will call IsUnsealed now sector=%+v, offset=%d, size=%d", sectorID, offset, length) - return rpn.pp.IsUnsealed(ctx, ref, storiface.UnpaddedByteIndex(offset), length) -} - // GetRetrievalPricingInput takes a set of candidate storage deals that can serve a retrieval request, // and returns an minimally populated PricingInput. This PricingInput should be enhanced // with more data, and passed to the pricing function to determine the final quoted price. @@ -175,37 +106,3 @@ func (rpn *retrievalProviderNode) GetRetrievalPricingInput(ctx context.Context, return resp, nil } - -func (rpn *retrievalProviderNode) sectorsStatus(ctx context.Context, sid abi.SectorNumber, showOnChainInfo bool) (api.SectorInfo, error) { - sInfo, err := rpn.secb.SectorsStatus(ctx, sid, false) - if err != nil { - return api.SectorInfo{}, err - } - - if !showOnChainInfo { - return sInfo, nil - } - - onChainInfo, err := rpn.full.StateSectorGetInfo(ctx, rpn.maddr, sid, types.EmptyTSK) - if err != nil { - return sInfo, err - } - if onChainInfo == nil { - return sInfo, nil - } - sInfo.SealProof = onChainInfo.SealProof - sInfo.Activation = onChainInfo.Activation - sInfo.Expiration = onChainInfo.Expiration - sInfo.DealWeight = onChainInfo.DealWeight - sInfo.VerifiedDealWeight = onChainInfo.VerifiedDealWeight - sInfo.InitialPledge = onChainInfo.InitialPledge - - ex, err := rpn.full.StateSectorExpiration(ctx, rpn.maddr, sid, types.EmptyTSK) - if err != nil { - return sInfo, nil - } - sInfo.OnTime = ex.OnTime - sInfo.Early = ex.Early - - return sInfo, nil -} diff --git a/markets/sectoraccessor/sectoraccessor.go b/markets/sectoraccessor/sectoraccessor.go new file mode 100644 index 000000000..33489a6fc --- /dev/null +++ b/markets/sectoraccessor/sectoraccessor.go @@ -0,0 +1,132 @@ +package sectoraccessor + +import ( + "context" + "io" + + "golang.org/x/xerrors" + + "github.com/filecoin-project/lotus/api" + "github.com/filecoin-project/lotus/api/v1api" + "github.com/filecoin-project/lotus/chain/types" + sectorstorage "github.com/filecoin-project/lotus/extern/sector-storage" + "github.com/filecoin-project/lotus/extern/sector-storage/storiface" + "github.com/filecoin-project/lotus/node/modules/dtypes" + "github.com/filecoin-project/lotus/storage/sectorblocks" + + "github.com/filecoin-project/go-address" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" + "github.com/filecoin-project/go-state-types/abi" + specstorage "github.com/filecoin-project/specs-storage/storage" + + "github.com/ipfs/go-cid" + logging "github.com/ipfs/go-log/v2" +) + +var log = logging.Logger("sectoraccessor") + +type sectorAccessor struct { + maddr address.Address + secb sectorblocks.SectorBuilder + pp sectorstorage.PieceProvider + full v1api.FullNode +} + +var _ sectoraccessor.SectorAccessor = (*sectorAccessor)(nil) + +func NewSectorAccessor(maddr dtypes.MinerAddress, secb sectorblocks.SectorBuilder, pp sectorstorage.PieceProvider, full v1api.FullNode) sectoraccessor.SectorAccessor { + return §orAccessor{address.Address(maddr), secb, pp, full} +} + +func (sa *sectorAccessor) UnsealSector(ctx context.Context, sectorID abi.SectorNumber, offset abi.UnpaddedPieceSize, length abi.UnpaddedPieceSize) (io.ReadCloser, error) { + log.Debugf("get sector %d, offset %d, length %d", sectorID, offset, length) + si, err := sa.sectorsStatus(ctx, sectorID, false) + if err != nil { + return nil, err + } + + mid, err := address.IDFromAddress(sa.maddr) + if err != nil { + return nil, err + } + + ref := specstorage.SectorRef{ + ID: abi.SectorID{ + Miner: abi.ActorID(mid), + Number: sectorID, + }, + ProofType: si.SealProof, + } + + var commD cid.Cid + if si.CommD != nil { + commD = *si.CommD + } + + // Get a reader for the piece, unsealing the piece if necessary + log.Debugf("read piece in sector %d, offset %d, length %d from miner %d", sectorID, offset, length, mid) + r, unsealed, err := sa.pp.ReadPiece(ctx, ref, storiface.UnpaddedByteIndex(offset), length, si.Ticket.Value, commD) + if err != nil { + return nil, xerrors.Errorf("failed to unseal piece from sector %d: %w", sectorID, err) + } + _ = unsealed // todo: use + + return r, nil +} + +func (sa *sectorAccessor) IsUnsealed(ctx context.Context, sectorID abi.SectorNumber, offset abi.UnpaddedPieceSize, length abi.UnpaddedPieceSize) (bool, error) { + si, err := sa.sectorsStatus(ctx, sectorID, true) + if err != nil { + return false, xerrors.Errorf("failed to get sector info: %w", err) + } + + mid, err := address.IDFromAddress(sa.maddr) + if err != nil { + return false, err + } + + ref := specstorage.SectorRef{ + ID: abi.SectorID{ + Miner: abi.ActorID(mid), + Number: sectorID, + }, + ProofType: si.SealProof, + } + + log.Debugf("will call IsUnsealed now sector=%+v, offset=%d, size=%d", sectorID, offset, length) + return sa.pp.IsUnsealed(ctx, ref, storiface.UnpaddedByteIndex(offset), length) +} + +func (sa *sectorAccessor) sectorsStatus(ctx context.Context, sid abi.SectorNumber, showOnChainInfo bool) (api.SectorInfo, error) { + sInfo, err := sa.secb.SectorsStatus(ctx, sid, false) + if err != nil { + return api.SectorInfo{}, err + } + + if !showOnChainInfo { + return sInfo, nil + } + + onChainInfo, err := sa.full.StateSectorGetInfo(ctx, sa.maddr, sid, types.EmptyTSK) + if err != nil { + return sInfo, err + } + if onChainInfo == nil { + return sInfo, nil + } + sInfo.SealProof = onChainInfo.SealProof + sInfo.Activation = onChainInfo.Activation + sInfo.Expiration = onChainInfo.Expiration + sInfo.DealWeight = onChainInfo.DealWeight + sInfo.VerifiedDealWeight = onChainInfo.VerifiedDealWeight + sInfo.InitialPledge = onChainInfo.InitialPledge + + ex, err := sa.full.StateSectorExpiration(ctx, sa.maddr, sid, types.EmptyTSK) + if err != nil { + return sInfo, nil + } + sInfo.OnTime = ex.OnTime + sInfo.Early = ex.Early + + return sInfo, nil +} diff --git a/markets/storageadapter/provider.go b/markets/storageadapter/provider.go index 19a2ef1e1..b899c0810 100644 --- a/markets/storageadapter/provider.go +++ b/markets/storageadapter/provider.go @@ -7,7 +7,6 @@ import ( "io" "time" - "github.com/filecoin-project/go-fil-markets/retrievalmarket" "github.com/ipfs/go-cid" logging "github.com/ipfs/go-log/v2" "go.uber.org/fx" @@ -58,12 +57,10 @@ type ProviderNodeAdapter struct { maxDealCollateralMultiplier uint64 dsMatcher *dealStateMatcher scMgr *SectorCommittedManager - - rpn retrievalmarket.RetrievalProviderNode } -func NewProviderNodeAdapter(fc *config.MinerFeeConfig, dc *config.DealmakingConfig) func(mctx helpers.MetricsCtx, lc fx.Lifecycle, dag dtypes.StagingDAG, secb *sectorblocks.SectorBlocks, full v1api.FullNode, dealPublisher *DealPublisher, rpn retrievalmarket.RetrievalProviderNode) storagemarket.StorageProviderNode { - return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, dag dtypes.StagingDAG, secb *sectorblocks.SectorBlocks, full v1api.FullNode, dealPublisher *DealPublisher, rpn retrievalmarket.RetrievalProviderNode) storagemarket.StorageProviderNode { +func NewProviderNodeAdapter(fc *config.MinerFeeConfig, dc *config.DealmakingConfig) func(mctx helpers.MetricsCtx, lc fx.Lifecycle, dag dtypes.StagingDAG, secb *sectorblocks.SectorBlocks, full v1api.FullNode, dealPublisher *DealPublisher) storagemarket.StorageProviderNode { + return func(mctx helpers.MetricsCtx, lc fx.Lifecycle, dag dtypes.StagingDAG, secb *sectorblocks.SectorBlocks, full v1api.FullNode, dealPublisher *DealPublisher) storagemarket.StorageProviderNode { ctx := helpers.LifecycleCtx(mctx, lc) ev := events.NewEvents(ctx, full) @@ -84,7 +81,6 @@ func NewProviderNodeAdapter(fc *config.MinerFeeConfig, dc *config.DealmakingConf na.maxDealCollateralMultiplier = dc.MaxProviderCollateralMultiplier } na.scMgr = NewSectorCommittedManager(ev, na, &apiWrapper{api: full}) - na.rpn = rpn return na } @@ -426,8 +422,4 @@ func (n *ProviderNodeAdapter) OnDealExpiredOrSlashed(ctx context.Context, dealID return nil } -func (n *ProviderNodeAdapter) IsUnsealed(ctx context.Context, sectorID abi.SectorNumber, offset abi.UnpaddedPieceSize, length abi.UnpaddedPieceSize) (bool, error) { - return n.rpn.IsUnsealed(ctx, sectorID, offset, length) -} - var _ storagemarket.StorageProviderNode = &ProviderNodeAdapter{} diff --git a/node/builder_miner.go b/node/builder_miner.go index 0bfcfdd25..d771fd4a3 100644 --- a/node/builder_miner.go +++ b/node/builder_miner.go @@ -4,6 +4,9 @@ import ( "errors" "time" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" + lotussectoraccessor "github.com/filecoin-project/lotus/markets/sectoraccessor" + "go.uber.org/fx" "golang.org/x/xerrors" @@ -152,6 +155,7 @@ func ConfigStorageMiner(c interface{}) Option { Override(DAGStoreKey, modules.DAGStore), // Markets (retrieval) + Override(new(sectoraccessor.SectorAccessor), lotussectoraccessor.NewSectorAccessor), Override(new(retrievalmarket.RetrievalProviderNode), retrievaladapter.NewRetrievalProviderNode), Override(new(rmnet.RetrievalMarketNetwork), modules.RetrievalNetwork), Override(new(retrievalmarket.RetrievalProvider), modules.RetrievalProvider), diff --git a/node/impl/client/client.go b/node/impl/client/client.go index 9dfb71019..a0ace9876 100644 --- a/node/impl/client/client.go +++ b/node/impl/client/client.go @@ -7,6 +7,7 @@ import ( "io" "math/rand" "os" + "path/filepath" "sort" "time" @@ -61,6 +62,7 @@ import ( "github.com/filecoin-project/lotus/node/impl/full" "github.com/filecoin-project/lotus/node/impl/paych" "github.com/filecoin-project/lotus/node/modules/dtypes" + "github.com/filecoin-project/lotus/node/repo" "github.com/filecoin-project/lotus/node/repo/importmgr" ) @@ -68,6 +70,7 @@ var DefaultHashFunction = uint64(mh.BLAKE2B_MIN + 31) // 8 days ~= SealDuration + PreCommit + MaxProveCommitDuration + 8 hour buffer const dealStartBufferHours uint64 = 8 * 24 +const DefaultDAGStoreDir = "dagstore" type API struct { fx.In @@ -88,6 +91,7 @@ type API struct { Host host.Host RetrievalStoreMgr dtypes.ClientRetrievalStoreManager + Repo repo.LockedRepo } func calcDealExpiration(minDuration uint64, md *dline.Info, startEpoch abi.ChainEpoch) abi.ChainEpoch { @@ -779,6 +783,13 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref } }) + tmpCarv2FilePath, err := a.getTmpCarV2FilePath() + if err != nil { + unsubscribe() + finish(xerrors.Errorf("Retrieve failed: %w", err)) + return + } + resp, err := a.Retrieval.Retrieve( ctx, order.Root, @@ -786,7 +797,8 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref order.Total, *order.MinerPeer, order.Client, - order.Miner) + order.Miner, + tmpCarv2FilePath) if err != nil { unsubscribe() @@ -887,6 +899,17 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref return } +// TODO: Come up with a better mechanism for creating the tmp CARv2 file path +func (a *API) getTmpCarV2FilePath() (string, error) { + carsPath := filepath.Join(a.Repo.Path(), DefaultDAGStoreDir, "retrieval-cars") + + if err := os.MkdirAll(carsPath, 0755); err != nil { + return "", xerrors.Errorf("failed to create dir") + } + + return filepath.Join(carsPath, fmt.Sprintf("%d.car", time.Now().UnixNano())), nil +} + func (a *API) ClientListRetrievals(ctx context.Context) ([]api.RetrievalInfo, error) { deals, err := a.Retrieval.ListDeals() if err != nil { diff --git a/node/impl/storminer.go b/node/impl/storminer.go index 88ad554c9..13049e4b1 100644 --- a/node/impl/storminer.go +++ b/node/impl/storminer.go @@ -10,6 +10,8 @@ import ( "strconv" "time" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" + "github.com/filecoin-project/dagstore" "github.com/filecoin-project/dagstore/shard" "github.com/filecoin-project/go-jsonrpc/auth" @@ -65,15 +67,15 @@ type StorageMinerAPI struct { RemoteStore *stores.Remote // Markets - PieceStore dtypes.ProviderPieceStore `optional:"true"` - StorageProvider storagemarket.StorageProvider `optional:"true"` - RetrievalProvider retrievalmarket.RetrievalProvider `optional:"true"` - RetrievalProviderNode retrievalmarket.RetrievalProviderNode `optional:"true"` - DataTransfer dtypes.ProviderDataTransfer `optional:"true"` - DealPublisher *storageadapter.DealPublisher `optional:"true"` - SectorBlocks *sectorblocks.SectorBlocks `optional:"true"` - Host host.Host `optional:"true"` - DAGStore *dagstore.DAGStore `optional:"true"` + PieceStore dtypes.ProviderPieceStore `optional:"true"` + StorageProvider storagemarket.StorageProvider `optional:"true"` + RetrievalProvider retrievalmarket.RetrievalProvider `optional:"true"` + SectorAccessor sectoraccessor.SectorAccessor `optional:"true"` + DataTransfer dtypes.ProviderDataTransfer `optional:"true"` + DealPublisher *storageadapter.DealPublisher `optional:"true"` + SectorBlocks *sectorblocks.SectorBlocks `optional:"true"` + Host host.Host `optional:"true"` + DAGStore *dagstore.DAGStore `optional:"true"` // Miner / storage Miner *storage.Miner `optional:"true"` @@ -630,8 +632,8 @@ func (sm *StorageMinerAPI) DagstoreInitializeAll(ctx context.Context, params api return nil, fmt.Errorf("dagstore not available on this node") } - if sm.RetrievalProviderNode == nil { - return nil, fmt.Errorf("retrieval provider node not available on this node") + if sm.SectorAccessor == nil { + return nil, fmt.Errorf("sector accessor not available on this node") } // prepare the thottler tokens. @@ -670,7 +672,7 @@ func (sm *StorageMinerAPI) DagstoreInitializeAll(ctx context.Context, params api var isUnsealed bool for _, d := range pi.Deals { - isUnsealed, err = sm.RetrievalProviderNode.IsUnsealed(ctx, d.SectorID, d.Offset.Unpadded(), d.Length.Unpadded()) + isUnsealed, err = sm.SectorAccessor.IsUnsealed(ctx, d.SectorID, d.Offset.Unpadded(), d.Length.Unpadded()) if err != nil { log.Warnw("DagstoreInitializeAll: failed to get unsealed status; skipping deal", "deal_id", d.DealID, "error", err) continue diff --git a/node/modules/client.go b/node/modules/client.go index b7757f8df..9e59183e0 100644 --- a/node/modules/client.go +++ b/node/modules/client.go @@ -182,16 +182,10 @@ func StorageClient(lc fx.Lifecycle, h host.Host, dataTransfer dtypes.ClientDataT func RetrievalClient(lc fx.Lifecycle, h host.Host, r repo.LockedRepo, dt dtypes.ClientDataTransfer, payAPI payapi.PaychAPI, resolver discovery.PeerResolver, ds dtypes.MetadataDS, chainAPI full.ChainAPI, stateAPI full.StateAPI, j journal.Journal) (retrievalmarket.RetrievalClient, error) { - carsPath := filepath.Join(r.Path(), DefaultDAGStoreDir, "retrieval-cars") - - if err := os.MkdirAll(carsPath, 0755); err != nil { - return nil, xerrors.Errorf("failed to create dir") - } - adapter := retrievaladapter.NewRetrievalClientNode(payAPI, chainAPI, stateAPI) network := rmnet.NewFromLibp2pHost(h) client, err := retrievalimpl.NewClient(network, - carsPath, dt, adapter, resolver, namespace.Wrap(ds, datastore.NewKey("/retrievals/client"))) + dt, adapter, resolver, namespace.Wrap(ds, datastore.NewKey("/retrievals/client"))) if err != nil { return nil, err } diff --git a/node/modules/storageminer.go b/node/modules/storageminer.go index 876ba7b27..17fe41883 100644 --- a/node/modules/storageminer.go +++ b/node/modules/storageminer.go @@ -42,6 +42,7 @@ import ( "github.com/filecoin-project/go-fil-markets/retrievalmarket" retrievalimpl "github.com/filecoin-project/go-fil-markets/retrievalmarket/impl" rmnet "github.com/filecoin-project/go-fil-markets/retrievalmarket/network" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" "github.com/filecoin-project/go-fil-markets/shared" "github.com/filecoin-project/go-fil-markets/storagemarket" storageimpl "github.com/filecoin-project/go-fil-markets/storagemarket/impl" @@ -670,6 +671,7 @@ func RetrievalPricingFunc(cfg config.DealmakingConfig) func(_ dtypes.ConsiderOnl func RetrievalProvider( maddr dtypes.MinerAddress, adapter retrievalmarket.RetrievalProviderNode, + sa sectoraccessor.SectorAccessor, netwk rmnet.RetrievalMarketNetwork, ds dtypes.MetadataDS, pieceStore dtypes.ProviderPieceStore, @@ -682,6 +684,7 @@ func RetrievalProvider( return retrievalimpl.NewProvider( address.Address(maddr), adapter, + sa, netwk, pieceStore, dagStore, diff --git a/node/modules/storageminer_dagstore.go b/node/modules/storageminer_dagstore.go index 9e6560290..d385eea5c 100644 --- a/node/modules/storageminer_dagstore.go +++ b/node/modules/storageminer_dagstore.go @@ -7,11 +7,11 @@ import ( "path/filepath" "strconv" - "github.com/filecoin-project/dagstore" "go.uber.org/fx" "golang.org/x/xerrors" - "github.com/filecoin-project/go-fil-markets/retrievalmarket" + "github.com/filecoin-project/dagstore" + "github.com/filecoin-project/go-fil-markets/sectoraccessor" mdagstore "github.com/filecoin-project/lotus/markets/dagstore" "github.com/filecoin-project/lotus/node/config" @@ -25,7 +25,7 @@ const ( ) // NewMinerAPI creates a new MinerAPI adaptor for the dagstore mounts. -func NewMinerAPI(lc fx.Lifecycle, r repo.LockedRepo, pieceStore dtypes.ProviderPieceStore, rpn retrievalmarket.RetrievalProviderNode) (mdagstore.MinerAPI, error) { +func NewMinerAPI(lc fx.Lifecycle, r repo.LockedRepo, pieceStore dtypes.ProviderPieceStore, sa sectoraccessor.SectorAccessor) (mdagstore.MinerAPI, error) { cfg, err := extractDAGStoreConfig(r) if err != nil { return nil, err @@ -40,7 +40,7 @@ func NewMinerAPI(lc fx.Lifecycle, r repo.LockedRepo, pieceStore dtypes.ProviderP } } - mountApi := mdagstore.NewMinerAPI(pieceStore, rpn, cfg.MaxConcurrencyStorageCalls) + mountApi := mdagstore.NewMinerAPI(pieceStore, sa, cfg.MaxConcurrencyStorageCalls) ready := make(chan error, 1) pieceStore.OnReady(func(err error) { ready <- err