lotus/api/test/deals.go

459 lines
11 KiB
Go
Raw Normal View History

2019-11-06 16:10:44 +00:00
package test
import (
2019-12-01 21:52:24 +00:00
"bytes"
2019-11-06 16:10:44 +00:00
"context"
"fmt"
2019-12-01 21:52:24 +00:00
"io/ioutil"
2019-11-06 16:10:44 +00:00
"math/rand"
2019-11-07 00:18:06 +00:00
"os"
2019-12-01 21:52:24 +00:00
"path/filepath"
"sync/atomic"
2019-11-06 16:10:44 +00:00
"testing"
"time"
2020-06-26 15:28:05 +00:00
"github.com/stretchr/testify/require"
2020-06-26 15:28:05 +00:00
"github.com/ipfs/go-cid"
files "github.com/ipfs/go-ipfs-files"
"github.com/ipld/go-car"
2019-11-06 19:44:28 +00:00
"github.com/filecoin-project/go-fil-markets/storagemarket"
2020-03-04 05:03:35 +00:00
"github.com/filecoin-project/lotus/api"
2019-11-30 23:17:50 +00:00
"github.com/filecoin-project/lotus/build"
2020-08-17 13:39:33 +00:00
sealing "github.com/filecoin-project/lotus/extern/storage-sealing"
dag "github.com/ipfs/go-merkledag"
dstest "github.com/ipfs/go-merkledag/test"
unixfile "github.com/ipfs/go-unixfs/file"
2019-11-06 16:10:44 +00:00
"github.com/filecoin-project/lotus/chain/types"
"github.com/filecoin-project/lotus/node/impl"
ipld "github.com/ipfs/go-ipld-format"
2019-11-06 16:10:44 +00:00
)
2020-07-08 18:35:55 +00:00
func TestDealFlow(t *testing.T, b APIBuilder, blocktime time.Duration, carExport, fastRet bool) {
2019-11-07 00:18:06 +00:00
2019-11-07 12:03:18 +00:00
ctx := context.Background()
n, sn := b(t, OneFull, OneMiner)
2019-11-06 16:10:44 +00:00
client := n[0].FullNode.(*impl.FullNodeAPI)
miner := sn[0]
addrinfo, err := client.NetAddrsListen(ctx)
if err != nil {
t.Fatal(err)
}
if err := miner.NetConnect(ctx, addrinfo); err != nil {
t.Fatal(err)
}
time.Sleep(time.Second)
mine := int64(1)
done := make(chan struct{})
go func() {
defer close(done)
for atomic.LoadInt64(&mine) == 1 {
time.Sleep(blocktime)
2020-07-28 16:16:46 +00:00
if err := sn[0].MineOne(ctx, MineNext); err != nil {
t.Error(err)
}
}
}()
2019-12-02 14:26:39 +00:00
MakeDeal(t, ctx, 6, client, miner, carExport, fastRet)
atomic.AddInt64(&mine, -1)
fmt.Println("shutting down mining")
<-done
}
func TestDoubleDealFlow(t *testing.T, b APIBuilder, blocktime time.Duration) {
ctx := context.Background()
n, sn := b(t, OneFull, OneMiner)
client := n[0].FullNode.(*impl.FullNodeAPI)
miner := sn[0]
addrinfo, err := client.NetAddrsListen(ctx)
2019-11-06 16:10:44 +00:00
if err != nil {
t.Fatal(err)
}
if err := miner.NetConnect(ctx, addrinfo); err != nil {
2019-11-06 16:10:44 +00:00
t.Fatal(err)
}
time.Sleep(time.Second)
2019-11-06 16:10:44 +00:00
mine := int64(1)
2019-11-07 00:18:06 +00:00
done := make(chan struct{})
2019-11-06 16:10:44 +00:00
go func() {
2019-11-07 00:18:06 +00:00
defer close(done)
for atomic.LoadInt64(&mine) == 1 {
time.Sleep(blocktime)
2020-07-28 16:16:46 +00:00
if err := sn[0].MineOne(ctx, MineNext); err != nil {
2019-12-05 05:14:19 +00:00
t.Error(err)
2019-11-06 19:44:28 +00:00
}
2019-11-06 16:10:44 +00:00
}
}()
2020-04-23 17:50:52 +00:00
MakeDeal(t, ctx, 6, client, miner, false, false)
MakeDeal(t, ctx, 7, client, miner, false, false)
atomic.AddInt64(&mine, -1)
fmt.Println("shutting down mining")
<-done
}
func MakeDeal(t *testing.T, ctx context.Context, rseed int, client api.FullNode, miner TestStorageNode, carExport, fastRet bool) {
res, data, err := CreateClientFile(ctx, client, rseed)
if err != nil {
t.Fatal(err)
}
fcid := res.Root
fmt.Println("FILE CID: ", fcid)
2020-07-08 18:35:55 +00:00
deal := startDeal(t, ctx, miner, client, fcid, fastRet)
2020-04-23 17:50:52 +00:00
// TODO: this sleep is only necessary because deals don't immediately get logged in the dealstore, we should fix this
time.Sleep(time.Second)
waitDealSealed(t, ctx, miner, client, deal, false)
2020-04-23 17:50:52 +00:00
// Retrieval
2020-07-09 16:29:57 +00:00
info, err := client.ClientGetDealInfo(ctx, *deal)
require.NoError(t, err)
2020-04-23 17:50:52 +00:00
testRetrieval(t, ctx, client, fcid, &info.PieceCID, carExport, data)
2020-04-23 17:50:52 +00:00
}
func CreateClientFile(ctx context.Context, client api.FullNode, rseed int) (*api.ImportRes, []byte, error) {
data := make([]byte, 1600)
rand.New(rand.NewSource(int64(rseed))).Read(data)
dir, err := ioutil.TempDir(os.TempDir(), "test-make-deal-")
if err != nil {
return nil, nil, err
}
path := filepath.Join(dir, "sourcefile.dat")
err = ioutil.WriteFile(path, data, 0644)
if err != nil {
return nil, nil, err
}
res, err := client.ClientImport(ctx, api.FileRef{Path: path})
if err != nil {
return nil, nil, err
}
return res, data, nil
}
2020-08-06 20:16:55 +00:00
func TestFastRetrievalDealFlow(t *testing.T, b APIBuilder, blocktime time.Duration) {
ctx := context.Background()
n, sn := b(t, OneFull, OneMiner)
2020-08-06 20:16:55 +00:00
client := n[0].FullNode.(*impl.FullNodeAPI)
miner := sn[0]
addrinfo, err := client.NetAddrsListen(ctx)
if err != nil {
t.Fatal(err)
}
if err := miner.NetConnect(ctx, addrinfo); err != nil {
t.Fatal(err)
}
time.Sleep(time.Second)
mine := int64(1)
done := make(chan struct{})
go func() {
defer close(done)
for atomic.LoadInt64(&mine) == 1 {
time.Sleep(blocktime)
if err := sn[0].MineOne(ctx, MineNext); err != nil {
t.Error(err)
}
}
}()
data := make([]byte, 1600)
rand.New(rand.NewSource(int64(8))).Read(data)
r := bytes.NewReader(data)
fcid, err := client.ClientImportLocal(ctx, r)
if err != nil {
t.Fatal(err)
}
fmt.Println("FILE CID: ", fcid)
deal := startDeal(t, ctx, miner, client, fcid, true)
waitDealPublished(t, ctx, miner, deal)
fmt.Println("deal published, retrieving")
// Retrieval
info, err := client.ClientGetDealInfo(ctx, *deal)
require.NoError(t, err)
testRetrieval(t, ctx, client, fcid, &info.PieceCID, false, data)
2020-08-06 20:16:55 +00:00
atomic.AddInt64(&mine, -1)
fmt.Println("shutting down mining")
<-done
}
func TestSenondDealRetrieval(t *testing.T, b APIBuilder, blocktime time.Duration) {
ctx := context.Background()
n, sn := b(t, OneFull, OneMiner)
client := n[0].FullNode.(*impl.FullNodeAPI)
miner := sn[0]
addrinfo, err := client.NetAddrsListen(ctx)
if err != nil {
t.Fatal(err)
}
if err := miner.NetConnect(ctx, addrinfo); err != nil {
t.Fatal(err)
}
time.Sleep(time.Second)
mine := int64(1)
done := make(chan struct{})
go func() {
defer close(done)
for atomic.LoadInt64(&mine) == 1 {
time.Sleep(blocktime)
2020-07-28 16:16:46 +00:00
if err := sn[0].MineOne(ctx, MineNext); err != nil {
t.Error(err)
}
}
}()
{
data1 := make([]byte, 800)
rand.New(rand.NewSource(int64(3))).Read(data1)
r := bytes.NewReader(data1)
fcid1, err := client.ClientImportLocal(ctx, r)
if err != nil {
t.Fatal(err)
}
data2 := make([]byte, 800)
rand.New(rand.NewSource(int64(9))).Read(data2)
r2 := bytes.NewReader(data2)
fcid2, err := client.ClientImportLocal(ctx, r2)
if err != nil {
t.Fatal(err)
}
deal1 := startDeal(t, ctx, miner, client, fcid1, true)
// TODO: this sleep is only necessary because deals don't immediately get logged in the dealstore, we should fix this
time.Sleep(time.Second)
waitDealSealed(t, ctx, miner, client, deal1, true)
deal2 := startDeal(t, ctx, miner, client, fcid2, true)
time.Sleep(time.Second)
waitDealSealed(t, ctx, miner, client, deal2, false)
// Retrieval
info, err := client.ClientGetDealInfo(ctx, *deal2)
require.NoError(t, err)
rf, _ := miner.SectorsRefs(ctx)
fmt.Printf("refs: %+v\n", rf)
testRetrieval(t, ctx, client, fcid2, &info.PieceCID, false, data2)
}
atomic.AddInt64(&mine, -1)
fmt.Println("shutting down mining")
<-done
}
func startDeal(t *testing.T, ctx context.Context, miner TestStorageNode, client api.FullNode, fcid cid.Cid, fastRet bool) *cid.Cid {
2020-04-23 17:50:52 +00:00
maddr, err := miner.ActorAddress(ctx)
if err != nil {
t.Fatal(err)
}
addr, err := client.WalletDefaultAddress(ctx)
if err != nil {
t.Fatal(err)
}
2020-03-04 05:03:35 +00:00
deal, err := client.ClientStartDeal(ctx, &api.StartDealParams{
Data: &storagemarket.DataRef{
TransferType: storagemarket.TTGraphsync,
Root: fcid,
},
2020-04-16 21:43:39 +00:00
Wallet: addr,
Miner: maddr,
EpochPrice: types.NewInt(1000000),
2020-07-28 17:51:47 +00:00
MinBlocksDuration: uint64(build.MinDealDuration),
2020-07-08 18:35:55 +00:00
FastRetrieval: fastRet,
2020-03-04 05:03:35 +00:00
})
2019-11-06 16:10:44 +00:00
if err != nil {
2020-02-23 15:50:36 +00:00
t.Fatalf("%+v", err)
2019-11-06 16:10:44 +00:00
}
2020-04-23 17:50:52 +00:00
return deal
}
2019-11-06 19:44:28 +00:00
func waitDealSealed(t *testing.T, ctx context.Context, miner TestStorageNode, client api.FullNode, deal *cid.Cid, noseal bool) {
loop:
2019-11-06 19:44:28 +00:00
for {
di, err := client.ClientGetDealInfo(ctx, *deal)
if err != nil {
t.Fatal(err)
}
2019-11-06 20:39:07 +00:00
switch di.State {
case storagemarket.StorageDealAwaitingPreCommit, storagemarket.StorageDealSealing:
if noseal {
return
}
2020-06-26 15:28:05 +00:00
startSealingWaiting(t, ctx, miner)
case storagemarket.StorageDealProposalRejected:
2019-11-06 20:39:07 +00:00
t.Fatal("deal rejected")
case storagemarket.StorageDealFailing:
2019-11-06 20:39:07 +00:00
t.Fatal("deal failed")
case storagemarket.StorageDealError:
2020-04-17 21:23:30 +00:00
t.Fatal("deal errored", di.Message)
case storagemarket.StorageDealActive:
2019-11-06 20:39:07 +00:00
fmt.Println("COMPLETE", di)
2019-11-07 00:18:06 +00:00
break loop
2019-11-06 20:39:07 +00:00
}
fmt.Println("Deal state: ", storagemarket.DealStates[di.State])
2019-11-06 19:44:28 +00:00
time.Sleep(time.Second / 2)
}
2020-04-23 17:50:52 +00:00
}
2019-11-07 00:18:06 +00:00
2020-08-06 20:16:55 +00:00
func waitDealPublished(t *testing.T, ctx context.Context, miner TestStorageNode, deal *cid.Cid) {
subCtx, cancel := context.WithCancel(ctx)
defer cancel()
updates, err := miner.MarketGetDealUpdates(subCtx)
2020-08-06 20:16:55 +00:00
if err != nil {
t.Fatal(err)
}
for {
select {
case <-ctx.Done():
t.Fatal("context timeout")
case di := <-updates:
if deal.Equals(di.ProposalCid) {
switch di.State {
case storagemarket.StorageDealProposalRejected:
t.Fatal("deal rejected")
case storagemarket.StorageDealFailing:
t.Fatal("deal failed")
case storagemarket.StorageDealError:
t.Fatal("deal errored", di.Message)
case storagemarket.StorageDealFinalizing, storagemarket.StorageDealAwaitingPreCommit, storagemarket.StorageDealSealing, storagemarket.StorageDealActive:
fmt.Println("COMPLETE", di)
return
}
fmt.Println("Deal state: ", storagemarket.DealStates[di.State])
2020-08-06 20:16:55 +00:00
}
}
}
}
2020-06-26 15:28:05 +00:00
func startSealingWaiting(t *testing.T, ctx context.Context, miner TestStorageNode) {
snums, err := miner.SectorsList(ctx)
require.NoError(t, err)
for _, snum := range snums {
si, err := miner.SectorsStatus(ctx, snum, false)
2020-06-26 15:28:05 +00:00
require.NoError(t, err)
2020-07-17 16:23:19 +00:00
t.Logf("Sector state: %s", si.State)
2020-06-26 15:28:05 +00:00
if si.State == api.SectorState(sealing.WaitDeals) {
require.NoError(t, miner.SectorStartSealing(ctx, snum))
}
}
}
func testRetrieval(t *testing.T, ctx context.Context, client api.FullNode, fcid cid.Cid, piece *cid.Cid, carExport bool, data []byte) {
2020-07-09 16:29:57 +00:00
offers, err := client.ClientFindData(ctx, fcid, piece)
2019-12-01 21:52:24 +00:00
if err != nil {
t.Fatal(err)
}
if len(offers) < 1 {
t.Fatal("no offers")
}
rpath, err := ioutil.TempDir("", "lotus-retrieve-test-")
if err != nil {
t.Fatal(err)
}
defer os.RemoveAll(rpath) //nolint:errcheck
2019-12-01 21:52:24 +00:00
caddr, err := client.WalletDefaultAddress(ctx)
if err != nil {
t.Fatal(err)
}
ref := &api.FileRef{
Path: filepath.Join(rpath, "ret"),
IsCAR: carExport,
}
updates, err := client.ClientRetrieveWithEvents(ctx, offers[0].Order(caddr), ref)
2020-10-03 00:45:15 +00:00
if err != nil {
t.Fatal(err)
}
2020-08-11 20:48:56 +00:00
for update := range updates {
if update.Err != "" {
2020-10-03 00:45:15 +00:00
t.Fatalf("retrieval failed: %s", update.Err)
2020-08-11 20:04:00 +00:00
}
2019-12-01 21:52:24 +00:00
}
rdata, err := ioutil.ReadFile(filepath.Join(rpath, "ret"))
if err != nil {
t.Fatal(err)
}
if carExport {
2020-04-23 17:50:52 +00:00
rdata = extractCarData(t, ctx, rdata, rpath)
}
2019-12-01 21:52:24 +00:00
if !bytes.Equal(rdata, data) {
t.Fatal("wrong data retrieved")
}
2020-04-23 17:50:52 +00:00
}
2019-12-01 21:52:24 +00:00
2020-04-23 17:50:52 +00:00
func extractCarData(t *testing.T, ctx context.Context, rdata []byte, rpath string) []byte {
bserv := dstest.Bserv()
ch, err := car.LoadCar(bserv.Blockstore(), bytes.NewReader(rdata))
if err != nil {
t.Fatal(err)
}
b, err := bserv.GetBlock(ctx, ch.Roots[0])
if err != nil {
t.Fatal(err)
}
nd, err := ipld.Decode(b)
if err != nil {
t.Fatal(err)
}
dserv := dag.NewDAGService(bserv)
fil, err := unixfile.NewUnixfsFile(ctx, dserv, nd)
if err != nil {
t.Fatal(err)
}
outPath := filepath.Join(rpath, "retLoadedCAR")
if err := files.WriteTo(fil, outPath); err != nil {
t.Fatal(err)
}
rdata, err = ioutil.ReadFile(outPath)
if err != nil {
t.Fatal(err)
}
return rdata
2019-11-06 16:10:44 +00:00
}