remove multistore from Lotus
This commit is contained in:
+65
-85
@@ -5,6 +5,7 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/rand"
|
||||
"os"
|
||||
"sort"
|
||||
"time"
|
||||
@@ -22,7 +23,6 @@ import (
|
||||
"github.com/ipfs/go-cid"
|
||||
offline "github.com/ipfs/go-ipfs-exchange-offline"
|
||||
files "github.com/ipfs/go-ipfs-files"
|
||||
ipld "github.com/ipfs/go-ipld-format"
|
||||
"github.com/ipfs/go-merkledag"
|
||||
unixfile "github.com/ipfs/go-unixfs/file"
|
||||
"github.com/ipld/go-car"
|
||||
@@ -46,7 +46,6 @@ import (
|
||||
"github.com/filecoin-project/go-fil-markets/shared"
|
||||
"github.com/filecoin-project/go-fil-markets/storagemarket"
|
||||
"github.com/filecoin-project/go-fil-markets/storagemarket/network"
|
||||
"github.com/filecoin-project/go-multistore"
|
||||
"github.com/filecoin-project/go-state-types/abi"
|
||||
"github.com/filecoin-project/specs-actors/v3/actors/builtin/market"
|
||||
|
||||
@@ -82,12 +81,10 @@ type API struct {
|
||||
Chain *store.ChainStore
|
||||
|
||||
Imports dtypes.ClientImportMgr
|
||||
Mds dtypes.ClientMultiDstore
|
||||
|
||||
CombinedBstore dtypes.ClientBlockstore // TODO: try to remove
|
||||
RetrievalStoreMgr dtypes.ClientRetrievalStoreManager
|
||||
DataTransfer dtypes.ClientDataTransfer
|
||||
Host host.Host
|
||||
CombinedBstore dtypes.ClientBlockstore // TODO: try to remove
|
||||
DataTransfer dtypes.ClientDataTransfer
|
||||
Host host.Host
|
||||
|
||||
// TODO How do we inject the Repo Path here ?
|
||||
}
|
||||
@@ -416,16 +413,29 @@ func (a *API) newDealInfoWithTransfer(transferCh *api.DataTransferChannel, v sto
|
||||
|
||||
func (a *API) ClientHasLocal(ctx context.Context, root cid.Cid) (bool, error) {
|
||||
// TODO: check if we have the ENTIRE dag
|
||||
|
||||
offExch := merkledag.NewDAGService(blockservice.New(a.Imports.Blockstore, offline.Exchange(a.Imports.Blockstore)))
|
||||
_, err := offExch.Get(ctx, root)
|
||||
if err == ipld.ErrNotFound {
|
||||
return false, nil
|
||||
}
|
||||
importIDs, err := a.imgr().List()
|
||||
if err != nil {
|
||||
return false, err
|
||||
return false, xerrors.Errorf("failed to list imports: %w", err)
|
||||
}
|
||||
return true, nil
|
||||
|
||||
for _, importID := range importIDs {
|
||||
info, err := a.imgr().Info(importID)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
if info.Labels[importmgr.LRootCid] == "" {
|
||||
continue
|
||||
}
|
||||
c, err := cid.Parse(info.Labels[importmgr.LRootCid])
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
if c.Equals(root) {
|
||||
return true, nil
|
||||
}
|
||||
}
|
||||
|
||||
return false, nil
|
||||
}
|
||||
|
||||
func (a *API) ClientFindData(ctx context.Context, root cid.Cid, piece *cid.Cid) ([]api.QueryOffer, error) {
|
||||
@@ -487,17 +497,11 @@ func (a *API) makeRetrievalQuery(ctx context.Context, rp rm.RetrievalPeer, paylo
|
||||
}
|
||||
|
||||
func (a *API) ClientImport(ctx context.Context, ref api.FileRef) (res *api.ImportRes, finalErr error) {
|
||||
id, st, err := a.imgr().NewStore()
|
||||
id, err := a.imgr().NewStore()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// we don't need the store any more after we return from here. clean it up completely.
|
||||
// we only need to retain the metadata related to this import which is identified by the storeID/importID.
|
||||
defer func() {
|
||||
_ = ((*multistore.MultiStore)(a.Mds)).Delete(id)
|
||||
}()
|
||||
|
||||
carV2File, err := a.imgr().NewTempFile(id)
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("failed to create temp CARv2 file: %w", err)
|
||||
@@ -516,7 +520,7 @@ func (a *API) ClientImport(ctx context.Context, ref api.FileRef) (res *api.Impor
|
||||
return nil, xerrors.Errorf("failed to import CAR file: %w", err)
|
||||
}
|
||||
} else {
|
||||
root, err = importNormalFileToCARv2(ctx, st, ref.Path, carV2File)
|
||||
root, err = a.importNormalFileToCARv2(ctx, id, ref.Path, carV2File)
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("failed to import normal file to CARv2: %w", err)
|
||||
}
|
||||
@@ -541,10 +545,10 @@ func (a *API) ClientImport(ctx context.Context, ref api.FileRef) (res *api.Impor
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (a *API) ClientRemoveImport(ctx context.Context, importID multistore.StoreID) error {
|
||||
func (a *API) ClientRemoveImport(ctx context.Context, importID uint64) error {
|
||||
info, err := a.imgr().Info(importID)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("failed to fetch multistore info: %w", err)
|
||||
return xerrors.Errorf("failed to fetch import info: %w", err)
|
||||
}
|
||||
|
||||
// remove the CARv2 file if we've created one.
|
||||
@@ -555,46 +559,33 @@ func (a *API) ClientRemoveImport(ctx context.Context, importID multistore.StoreI
|
||||
return a.imgr().Remove(importID)
|
||||
}
|
||||
|
||||
// FIXME
|
||||
func (a *API) ClientImportLocal(ctx context.Context, f io.Reader) (c cid.Cid, finalErr error) {
|
||||
id, st, err := a.imgr().NewStore()
|
||||
func (a *API) ClientImportLocal(ctx context.Context, r io.Reader) (cid.Cid, error) {
|
||||
// write payload to temp file
|
||||
tmpPath, err := a.imgr().NewTempFile(rand.Uint64())
|
||||
if err != nil {
|
||||
return cid.Undef, err
|
||||
}
|
||||
// we don't need the store any more after we return from here. clean it up completely.
|
||||
// we only need to retain the metadata related to this import which is identified by the storeID/importID.
|
||||
defer func() {
|
||||
_ = ((*multistore.MultiStore)(a.Mds)).Delete(id)
|
||||
}()
|
||||
|
||||
carV2File, err := a.imgr().NewTempFile(id)
|
||||
defer os.Remove(tmpPath) //nolint:errcheck
|
||||
tmpF, err := os.Open(tmpPath)
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to create temp CARv2 file: %w", err)
|
||||
return cid.Undef, err
|
||||
}
|
||||
// make sure to remove the CARv2 file if anything goes wrong from here on.
|
||||
defer func() {
|
||||
if finalErr != nil {
|
||||
_ = os.Remove(carV2File)
|
||||
}
|
||||
}()
|
||||
|
||||
// FIXME
|
||||
root, err := importNormalFileToCARv2(ctx, st, "", carV2File)
|
||||
if err != nil {
|
||||
return root, xerrors.Errorf("failed to import to CARv2 file: %w", err)
|
||||
defer tmpF.Close() //nolint:errcheck
|
||||
if _, err := io.Copy(tmpF, r); err != nil {
|
||||
return cid.Undef, err
|
||||
}
|
||||
|
||||
if err := a.imgr().AddLabel(id, importmgr.LSource, "import-local"); err != nil {
|
||||
return cid.Cid{}, err
|
||||
}
|
||||
if err := a.imgr().AddLabel(id, importmgr.LRootCid, root.String()); err != nil {
|
||||
return cid.Cid{}, err
|
||||
}
|
||||
if err := a.imgr().AddLabel(id, importmgr.LCARv2FilePath, carV2File); err != nil {
|
||||
if err := tmpF.Close(); err != nil {
|
||||
return cid.Undef, err
|
||||
}
|
||||
|
||||
return root, nil
|
||||
res, err := a.ClientImport(ctx, api.FileRef{
|
||||
Path: tmpPath,
|
||||
IsCAR: false,
|
||||
})
|
||||
if err != nil {
|
||||
return cid.Undef, err
|
||||
}
|
||||
return res.Root, nil
|
||||
}
|
||||
|
||||
func (a *API) ClientListImports(ctx context.Context) ([]api.Import, error) {
|
||||
@@ -818,8 +809,8 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref
|
||||
}
|
||||
|
||||
carV2FilePath = resp.CarFilePath
|
||||
// remove the temp CARv2 fil when retrieval is complete
|
||||
defer os.Remove(carV2FilePath)
|
||||
// remove the temp CARv2 file when retrieval is complete
|
||||
defer os.Remove(carV2FilePath) //nolint:errcheck
|
||||
} else {
|
||||
carV2FilePath = order.LocalCARV2FilePath
|
||||
}
|
||||
@@ -837,7 +828,7 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref
|
||||
finish(err)
|
||||
return
|
||||
}
|
||||
defer carv2Reader.Close()
|
||||
defer carv2Reader.Close() //nolint:errcheck
|
||||
if _, err := io.Copy(f, carv2Reader.CarV1Reader()); err != nil {
|
||||
finish(err)
|
||||
return
|
||||
@@ -852,7 +843,7 @@ func (a *API) clientRetrieve(ctx context.Context, order api.RetrievalOrder, ref
|
||||
finish(err)
|
||||
return
|
||||
}
|
||||
defer rw.Close()
|
||||
defer rw.Close() //nolint:errcheck
|
||||
bsvc := blockservice.New(rw, offline.Exchange(rw))
|
||||
dag := merkledag.NewDAGService(bsvc)
|
||||
|
||||
@@ -947,19 +938,6 @@ func (a *API) newRetrievalInfo(ctx context.Context, v rm.ClientDealState) api.Re
|
||||
return a.newRetrievalInfoWithTransfer(transferCh, v)
|
||||
}
|
||||
|
||||
type multiStoreRetrievalStore struct {
|
||||
storeID multistore.StoreID
|
||||
store *multistore.Store
|
||||
}
|
||||
|
||||
func (mrs *multiStoreRetrievalStore) StoreID() *multistore.StoreID {
|
||||
return &mrs.storeID
|
||||
}
|
||||
|
||||
func (mrs *multiStoreRetrievalStore) DAGService() ipld.DAGService {
|
||||
return mrs.store.DAG
|
||||
}
|
||||
|
||||
func (a *API) ClientQueryAsk(ctx context.Context, p peer.ID, miner address.Address) (*storagemarket.StorageAsk, error) {
|
||||
mi, err := a.StateMinerInfo(ctx, miner, types.EmptyTSK)
|
||||
if err != nil {
|
||||
@@ -1065,28 +1043,30 @@ func (a *API) ClientDealPieceCID(ctx context.Context, root cid.Cid) (api.DataCID
|
||||
}
|
||||
|
||||
func (a *API) ClientGenCar(ctx context.Context, ref api.FileRef, outputPath string) error {
|
||||
id, st, err := a.imgr().NewStore()
|
||||
id := rand.Uint64()
|
||||
tmpCARv2File, err := a.imgr().NewTempFile(id)
|
||||
if err != nil {
|
||||
return err
|
||||
return xerrors.Errorf("failed to create temp file: %w", err)
|
||||
}
|
||||
defer os.Remove(tmpCARv2File) //nolint:errcheck
|
||||
|
||||
defer func() {
|
||||
// Clean up the store as we don't need it anymore.
|
||||
_ = a.imgr().Remove(id)
|
||||
}()
|
||||
|
||||
c, err := importNormalFileToUnixfsDAG(ctx, ref.Path, st.DAG)
|
||||
root, err := a.importNormalFileToCARv2(ctx, id, ref.Path, tmpCARv2File)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("failed to import file to store: %w", err)
|
||||
return xerrors.Errorf("failed to import normal file to CARv2")
|
||||
}
|
||||
|
||||
// generate a deterministic CARv1 payload from the UnixFS DAG by doing an IPLD
|
||||
// traversal over the Unixfs DAGs in the blockstore using the "all selector" i.e the entire DAG selector.
|
||||
// traversal over the Unixfs DAG in the CARv2 file using the "all selector" i.e the entire DAG selector.
|
||||
rdOnly, err := blockstore.OpenReadOnly(tmpCARv2File, true)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("failed to open read only CARv2 blockstore: %w", err)
|
||||
}
|
||||
defer rdOnly.Close() //nolint:errcheck
|
||||
|
||||
ssb := builder.NewSelectorSpecBuilder(basicnode.Prototype.Any)
|
||||
allSelector := ssb.ExploreRecursive(selector.RecursionLimitNone(),
|
||||
ssb.ExploreAll(ssb.ExploreRecursiveEdge())).Node()
|
||||
|
||||
sc := car.NewSelectiveCar(ctx, st.Bstore, []car.Dag{{Root: c, Selector: allSelector}})
|
||||
sc := car.NewSelectiveCar(ctx, rdOnly, []car.Dag{{Root: root, Selector: allSelector}})
|
||||
f, err := os.Create(outputPath)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
+29
-27
@@ -3,11 +3,9 @@ package client
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
|
||||
"github.com/filecoin-project/go-multistore"
|
||||
"github.com/filecoin-project/lotus/build"
|
||||
"github.com/ipfs/go-blockservice"
|
||||
"github.com/ipfs/go-cid"
|
||||
@@ -25,43 +23,47 @@ import (
|
||||
"golang.org/x/xerrors"
|
||||
)
|
||||
|
||||
// importNormalFileToCARv2 imports the client's normal file to a CARv2 file.
|
||||
// It first generates a Unixfs DAG using the given store to store the resulting blocks and gets the root cid of the Unixfs DAG.
|
||||
// It then writes out the Unixfs DAG to a CARv2 file by generating a Unixfs DAG again using a CARv2 read-write blockstore as the backing store
|
||||
// and then finalizing the CARv2 read-write blockstore to get the backing CARv2 file.
|
||||
func importNormalFileToCARv2(ctx context.Context, st *multistore.Store, inputFilePath string, outputCARv2Path string) (c cid.Cid, finalErr error) {
|
||||
// create the UnixFS DAG and import the file to store to get the root.
|
||||
root, err := importNormalFileToUnixfsDAG(ctx, inputFilePath, st.DAG)
|
||||
// importNormalFileToCARv2 transforms the client's "normal file" to a Unixfs IPLD DAG and writes out the DAG to a CARv2 file at the given output path.
|
||||
func (a *API) importNormalFileToCARv2(ctx context.Context, importID uint64, inputFilePath string, outputCARv2Path string) (c cid.Cid, finalErr error) {
|
||||
|
||||
// TODO: We've currently put in a hack to create the Unixfs DAG as a CARv2 without using Badger.
|
||||
// We first transform the Unixfs DAG to a rootless CARv2 file as CARv2 doesen't allow streaming writes without specifying the root upfront and we
|
||||
// don't have the root till the Unixfs DAG is created.
|
||||
//
|
||||
// In the second pass, we create a CARv2 file with the root present using the root node we get in the above step.
|
||||
// This hack should be fixed when CARv2 allows specifying the root AFTER finishing the CARv2 streaming write.
|
||||
tmpCARv2Path, err := a.imgr().NewTempFile(importID)
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to create temp CARv2 file: %w", err)
|
||||
}
|
||||
defer os.Remove(tmpCARv2Path) //nolint:errcheck
|
||||
|
||||
tempCARv2Store, err := blockstore.NewReadWrite(tmpCARv2Path, []cid.Cid{})
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to create rootless temp CARv2 rw store: %w", err)
|
||||
}
|
||||
defer tempCARv2Store.Finalize() //nolint:errcheck
|
||||
bsvc := blockservice.New(tempCARv2Store, offline.Exchange(tempCARv2Store))
|
||||
|
||||
// ---- First Pass --- Write out the UnixFS DAG to a rootless CARv2 file by instantiating a read-write CARv2 blockstore without the root.
|
||||
root, err := importNormalFileToUnixfsDAG(ctx, inputFilePath, merkledag.NewDAGService(bsvc))
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to import file to store: %w", err)
|
||||
}
|
||||
|
||||
//---
|
||||
// transform the file to a CARv2 file by writing out a Unixfs DAG via the CARv2 read-write blockstore.
|
||||
//------ Second Pass --- Now that we have the root of the Unixfs DAG -> write out the Unixfs DAG to a CARv2 file with the root present by using a read-write CARv2 blockstore.
|
||||
rw, err := blockstore.NewReadWrite(outputCARv2Path, []cid.Cid{root})
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to create a CARv2 read-write blockstore: %w", err)
|
||||
}
|
||||
defer rw.Finalize() //nolint:errcheck
|
||||
|
||||
// make sure to call finalize on the CARv2 read-write blockstore to ensure that the blockstore flushes out a valid CARv2 file
|
||||
// and releases the file handle it acquires.
|
||||
defer func() {
|
||||
err := rw.Finalize()
|
||||
if finalErr != nil {
|
||||
finalErr = xerrors.Errorf("failed to import file to CARv2, err=%w", finalErr)
|
||||
} else {
|
||||
finalErr = err
|
||||
}
|
||||
}()
|
||||
|
||||
bsvc := blockservice.New(rw, offline.Exchange(rw))
|
||||
bsvc = blockservice.New(rw, offline.Exchange(rw))
|
||||
root2, err := importNormalFileToUnixfsDAG(ctx, inputFilePath, merkledag.NewDAGService(bsvc))
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to create Unixfs DAG with CARv2 blockstore: %w", err)
|
||||
}
|
||||
|
||||
fmt.Printf("\n root1 is %s and root2 is %s", root, root2)
|
||||
|
||||
if root != root2 {
|
||||
return cid.Undef, xerrors.New("roots do not match")
|
||||
}
|
||||
@@ -75,7 +77,7 @@ func importNormalFileToUnixfsDAG(ctx context.Context, inputFilePath string, dag
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to open input file: %w", err)
|
||||
}
|
||||
defer f.Close()
|
||||
defer f.Close() //nolint:errcheck
|
||||
|
||||
stat, err := f.Stat()
|
||||
if err != nil {
|
||||
@@ -147,7 +149,7 @@ func transformCarToCARv2(inputCARPath string, outputCARv2Path string) (root cid.
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to open output CARv2 file: %w", err)
|
||||
}
|
||||
defer outF.Close()
|
||||
defer outF.Close() //nolint:errcheck
|
||||
_, err = io.Copy(outF, inputF)
|
||||
if err != nil {
|
||||
return cid.Undef, xerrors.Errorf("failed to copy CARv2 file: %w", err)
|
||||
|
||||
Reference in New Issue
Block a user