Skip to content

Commit 9c642ba

Browse files
authored
Merge pull request #807 from textileio/jsign/stagecid
Eagerly fetch data in hot-storage
2 parents 1726c2a + 90a20cf commit 9c642ba

12 files changed

Lines changed: 139 additions & 96 deletions

File tree

api/client/utils_test.go

Lines changed: 26 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -36,31 +36,32 @@ func defaultServerConfig(t *testing.T) server.Config {
3636

3737
grpcMaddr := util.MustParseAddr(grpcHostAddress)
3838
conf := server.Config{
39-
WalletInitialFunds: *big.NewInt(int64(4000000000)),
40-
IpfsAPIAddr: ipfsAddr,
41-
LotusAddress: devnetAddr,
42-
LotusAuthToken: "",
43-
LotusMasterAddr: "",
44-
LotusConnectionRetries: 5,
45-
Devnet: true,
46-
GrpcHostNetwork: grpcHostNetwork,
47-
GrpcHostAddress: grpcMaddr,
48-
GrpcWebProxyAddress: grpcWebProxyAddress,
49-
RepoPath: repoPath,
50-
GatewayHostAddr: gatewayHostAddr,
51-
IndexRawJSONHostAddr: indexRawJSONHostAddr,
52-
MaxMindDBFolder: "../../iplocation/maxmind",
53-
MinerSelector: "reputation",
54-
FFSDealFinalityTimeout: time.Minute * 30,
55-
FFSMaxParallelDealPreparing: 1,
56-
FFSGCAutomaticGCInterval: 0,
57-
DealWatchPollDuration: time.Second * 15,
58-
SchedMaxParallel: 10,
59-
AskIndexQueryAskTimeout: time.Second * 3,
60-
AskIndexRefreshInterval: time.Second * 3,
61-
AskIndexRefreshOnStart: true,
62-
AskindexMaxParallel: 2,
63-
IndexMinersRefreshOnStart: false,
39+
WalletInitialFunds: *big.NewInt(int64(4000000000)),
40+
IpfsAPIAddr: ipfsAddr,
41+
LotusAddress: devnetAddr,
42+
LotusAuthToken: "",
43+
LotusMasterAddr: "",
44+
LotusConnectionRetries: 5,
45+
Devnet: true,
46+
GrpcHostNetwork: grpcHostNetwork,
47+
GrpcHostAddress: grpcMaddr,
48+
GrpcWebProxyAddress: grpcWebProxyAddress,
49+
RepoPath: repoPath,
50+
GatewayHostAddr: gatewayHostAddr,
51+
IndexRawJSONHostAddr: indexRawJSONHostAddr,
52+
MaxMindDBFolder: "../../iplocation/maxmind",
53+
MinerSelector: "reputation",
54+
FFSDealFinalityTimeout: time.Minute * 30,
55+
FFSMaxParallelDealPreparing: 1,
56+
FFSGCAutomaticGCInterval: 0,
57+
FFSRetrievalNextEventTimeout: time.Hour,
58+
DealWatchPollDuration: time.Second * 15,
59+
SchedMaxParallel: 10,
60+
AskIndexQueryAskTimeout: time.Second * 3,
61+
AskIndexRefreshInterval: time.Second * 3,
62+
AskIndexRefreshOnStart: true,
63+
AskindexMaxParallel: 2,
64+
IndexMinersRefreshOnStart: false,
6465
}
6566
return conf
6667
}

api/server/server.go

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -126,19 +126,20 @@ type Config struct {
126126
MongoURI string
127127
MongoDB string
128128

129-
FFSAdminToken string
130-
FFSUseMasterAddr bool
131-
FFSDealFinalityTimeout time.Duration
132-
FFSMinimumPieceSize uint64
133-
FFSMaxParallelDealPreparing int
134-
FFSGCAutomaticGCInterval time.Duration
135-
FFSGCStageGracePeriod time.Duration
136-
SchedMaxParallel int
137-
MinerSelector string
138-
MinerSelectorParams string
139-
DealWatchPollDuration time.Duration
140-
AutocreateMasterAddr bool
141-
WalletInitialFunds big.Int
129+
FFSAdminToken string
130+
FFSUseMasterAddr bool
131+
FFSDealFinalityTimeout time.Duration
132+
FFSMinimumPieceSize uint64
133+
FFSRetrievalNextEventTimeout time.Duration
134+
FFSMaxParallelDealPreparing int
135+
FFSGCAutomaticGCInterval time.Duration
136+
FFSGCStageGracePeriod time.Duration
137+
SchedMaxParallel int
138+
MinerSelector string
139+
MinerSelectorParams string
140+
DealWatchPollDuration time.Duration
141+
AutocreateMasterAddr bool
142+
WalletInitialFunds big.Int
142143

143144
AskIndexQueryAskTimeout time.Duration
144145
AskindexMaxParallel int
@@ -261,7 +262,7 @@ func NewServer(conf Config) (*Server, error) {
261262
if conf.Devnet {
262263
conf.FFSMinimumPieceSize = 0
263264
}
264-
cs := filcold.New(ms, dm, wm, ipfs, chain, l, lsm, conf.FFSMinimumPieceSize, conf.FFSMaxParallelDealPreparing)
265+
cs := filcold.New(ms, dm, wm, ipfs, chain, l, lsm, conf.FFSMinimumPieceSize, conf.FFSMaxParallelDealPreparing, conf.FFSRetrievalNextEventTimeout)
265266
hs, err := coreipfs.New(txndstr.Wrap(ds, "ffs/coreipfs"), ipfs, l)
266267
if err != nil {
267268
return nil, fmt.Errorf("creating coreipfs: %s", err)

cmd/powd/main.go

Lines changed: 15 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,7 @@ func configFromFlags() (server.Config, error) {
130130
ffsSchedMaxParallel := config.GetInt("ffsschedmaxparallel")
131131
ffsDealWatchFinalityTimeout := time.Minute * time.Duration(config.GetInt("ffsdealfinalitytimeout"))
132132
ffsMinimumPieceSize := config.GetUint64("ffsminimumpiecesize")
133+
ffsRetrievalNextEventTimeout := config.GetDuration("ffsretrievalnexteventtimeout")
133134
ffsMaxParallelDealPreparing := config.GetInt("ffsmaxparalleldealpreparing")
134135
ffsGCInterval := time.Minute * time.Duration(config.GetInt("ffsgcinterval"))
135136
ffsGCStagedGracePeriod := time.Minute * time.Duration(config.GetInt("ffsgcstagedgraceperiod"))
@@ -166,18 +167,19 @@ func configFromFlags() (server.Config, error) {
166167
MongoURI: mongoURI,
167168
MongoDB: mongoDB,
168169

169-
FFSAdminToken: ffsAdminToken,
170-
FFSUseMasterAddr: ffsUseMasterAddr,
171-
FFSDealFinalityTimeout: ffsDealWatchFinalityTimeout,
172-
FFSMinimumPieceSize: ffsMinimumPieceSize,
173-
FFSMaxParallelDealPreparing: ffsMaxParallelDealPreparing,
174-
FFSGCAutomaticGCInterval: ffsGCInterval,
175-
FFSGCStageGracePeriod: ffsGCStagedGracePeriod,
176-
AutocreateMasterAddr: autocreateMasterAddr,
177-
MinerSelector: minerSelector,
178-
MinerSelectorParams: minerSelectorParams,
179-
SchedMaxParallel: ffsSchedMaxParallel,
180-
DealWatchPollDuration: dealWatchPollDuration,
170+
FFSAdminToken: ffsAdminToken,
171+
FFSUseMasterAddr: ffsUseMasterAddr,
172+
FFSDealFinalityTimeout: ffsDealWatchFinalityTimeout,
173+
FFSMinimumPieceSize: ffsMinimumPieceSize,
174+
FFSRetrievalNextEventTimeout: ffsRetrievalNextEventTimeout,
175+
FFSMaxParallelDealPreparing: ffsMaxParallelDealPreparing,
176+
FFSGCAutomaticGCInterval: ffsGCInterval,
177+
FFSGCStageGracePeriod: ffsGCStagedGracePeriod,
178+
AutocreateMasterAddr: autocreateMasterAddr,
179+
MinerSelector: minerSelector,
180+
MinerSelectorParams: minerSelectorParams,
181+
SchedMaxParallel: ffsSchedMaxParallel,
182+
DealWatchPollDuration: dealWatchPollDuration,
181183

182184
AskIndexQueryAskTimeout: askIndexQueryAskTimeout,
183185
AskIndexRefreshInterval: askIndexRefreshInterval,
@@ -387,6 +389,7 @@ func setupFlags() error {
387389
pflag.String("ffsminerselector", "reputation", "Miner selector to be used by FFS: 'sr2', 'reputation'.")
388390
pflag.String("ffsminerselectorparams", "", "Miner selector configuration parameter, depends on --ffsminerselector.")
389391
pflag.String("ffsminimumpiecesize", "67108864", "Minimum piece size in bytes allowed to be stored in Filecoin.")
392+
pflag.Duration("ffsretrievalnexteventtimeout", time.Hour, "Maximum amount of time to wait for the next retrieval event before erroring it.")
390393
pflag.String("ffsschedmaxparallel", "1000", "Maximum amount of Jobs executed in parallel.")
391394
pflag.String("ffsdealfinalitytimeout", "4320", "Deadline in minutes in which a deal must prove liveness changing status before considered abandoned.")
392395
pflag.String("ffsmaxparalleldealpreparing", "2", "Max parallel deal preparing tasks.")

deals/module/records.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -411,7 +411,7 @@ func (m *Module) recordRetrieval(addr string, offer api.QueryOffer, bytesReceive
411411
RootCid: offer.Root,
412412
Size: offer.Size,
413413
MinPrice: offer.MinPrice.Uint64(),
414-
Miner: offer.Miner.String(),
414+
Miner: offer.MinerPeer.Address.String(),
415415
MinerPeerID: offer.MinerPeer.ID.String(),
416416
PaymentInterval: offer.PaymentInterval,
417417
PaymentIntervalIncrease: offer.PaymentIntervalIncrease,

deals/module/retrieve.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -123,7 +123,7 @@ func (m *Module) retrieve(ctx context.Context, lapi *apistruct.FullNodeStruct, l
123123
break Loop
124124
}
125125
if e.Err != "" {
126-
log.Infof("in progress retrieval errored: %s", err)
126+
log.Infof("in progress retrieval errored: %s", e.Err)
127127
errMsg = e.Err
128128
}
129129
if dtStart.IsZero() && e.Event == retrievalmarket.ClientEventBlocksReceived {
@@ -146,23 +146,23 @@ func (m *Module) retrieve(ctx context.Context, lapi *apistruct.FullNodeStruct, l
146146
// payment channel creation. This isn't ideal, but
147147
// it's better than missing the data.
148148
// We WARN just to signal this might be happening.
149-
if dtStart.IsZero() {
149+
if dtStart.IsZero() && errMsg == "" {
150150
dtStart = retrievalStartTime
151151
log.Warnf("retrieval data-transfer start fallback to retrieval start")
152152
}
153153
// This is a fallback to not receiving an expected
154154
// event in the retrieval. We just fallback to Now(),
155155
// which should always be pretty close to the real
156156
// event. We WARN just to signal this is happening.
157-
if dtEnd.IsZero() {
157+
if dtEnd.IsZero() && errMsg == "" {
158158
dtEnd = time.Now()
159159
log.Warnf("retrieval data-transfer end fallback to retrieval end")
160160
}
161161
m.recordRetrieval(waddr, o, bytesReceived, dtStart, dtEnd, errMsg)
162162
}
163163
}()
164164

165-
return o.Miner.String(), out, nil
165+
return o.MinerPeer.Address.String(), out, nil
166166
}
167167

168168
func getRetrievalOffers(ctx context.Context, lapi *apistruct.FullNodeStruct, payloadCid cid.Cid, pieceCid *cid.Cid, miners []string) []api.QueryOffer {

ffs/coreipfs/coreipfs.go

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -57,13 +57,12 @@ func New(ds datastore.TxnDatastore, ipfs iface.CoreAPI, l ffs.JobLogger) (*CoreI
5757

5858
// Stage adds the data of io.Reader in the storage, and creates a stage-pin on the resulting cid.
5959
func (ci *CoreIpfs) Stage(ctx context.Context, iid ffs.APIID, r io.Reader) (cid.Cid, error) {
60-
ci.lock.Lock()
61-
defer ci.lock.Unlock()
62-
6360
p, err := ci.ipfs.Unixfs().Add(ctx, ipfsfiles.NewReaderFile(r), options.Unixfs.Pin(true))
6461
if err != nil {
6562
return cid.Undef, fmt.Errorf("adding data to ipfs: %s", err)
6663
}
64+
ci.lock.Lock()
65+
defer ci.lock.Unlock()
6766

6867
if err := ci.ps.AddStaged(iid, p.Cid()); err != nil {
6968
return cid.Undef, fmt.Errorf("saving new pin in pinstore: %s", err)
@@ -72,8 +71,11 @@ func (ci *CoreIpfs) Stage(ctx context.Context, iid ffs.APIID, r io.Reader) (cid.
7271
return p.Cid(), nil
7372
}
7473

75-
// StageCid stage-pin a Cid.
74+
// StageCid pull the Cid data and stage-pin it.
7675
func (ci *CoreIpfs) StageCid(ctx context.Context, iid ffs.APIID, c cid.Cid) error {
76+
if err := ci.ipfs.Pin().Add(ctx, path.IpfsPath(c), options.Pin.Recursive(true)); err != nil {
77+
return fmt.Errorf("adding data to ipfs: %s", err)
78+
}
7779
ci.lock.Lock()
7880
defer ci.lock.Unlock()
7981

@@ -263,7 +265,7 @@ Loop:
263265

264266
// Skip Cids that are excluded.
265267
if _, ok := excludeMap[stagedPin.Cid]; ok {
266-
log.Infof("skipping staged cid %s since it's in exclusion list", stagedPin)
268+
log.Infof("skipping staged cid %s since it's in exclusion list", stagedPin.Cid)
267269
continue Loop
268270
}
269271
// A Cid is only safe to GC if all existing stage-pin are older than

ffs/filcold/filcold.go

Lines changed: 53 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212
"github.com/filecoin-project/go-fil-markets/storagemarket"
1313
"github.com/filecoin-project/go-state-types/abi"
1414
"github.com/filecoin-project/lotus/api"
15+
marketevents "github.com/filecoin-project/lotus/markets/loggers"
1516
"github.com/ipfs/go-cid"
1617
logger "github.com/ipfs/go-log/v2"
1718
iface "github.com/ipfs/interface-go-ipfs-core"
@@ -35,15 +36,16 @@ var (
3536
// FilCold is a ColdStorage implementation which saves data in the Filecoin network.
3637
// It assumes the underlying Filecoin client has access to an IPFS node where data is stored.
3738
type FilCold struct {
38-
ms ffs.MinerSelector
39-
dm *dealsModule.Module
40-
wm wallet.Module
41-
ipfs iface.CoreAPI
42-
chain FilChain
43-
l ffs.JobLogger
44-
lsm *lotus.SyncMonitor
45-
minPieceSize uint64
46-
semaphDealPrep chan struct{}
39+
ms ffs.MinerSelector
40+
dm *dealsModule.Module
41+
wm wallet.Module
42+
ipfs iface.CoreAPI
43+
chain FilChain
44+
l ffs.JobLogger
45+
lsm *lotus.SyncMonitor
46+
minPieceSize uint64
47+
retrNextEventTimeout time.Duration
48+
semaphDealPrep chan struct{}
4749
}
4850

4951
var _ ffs.ColdStorage = (*FilCold)(nil)
@@ -54,17 +56,18 @@ type FilChain interface {
5456
}
5557

5658
// New returns a new FilCold instance.
57-
func New(ms ffs.MinerSelector, dm *dealsModule.Module, wm wallet.Module, ipfs iface.CoreAPI, chain FilChain, l ffs.JobLogger, lsm *lotus.SyncMonitor, minPieceSize uint64, maxParallelDealPreparing int) *FilCold {
59+
func New(ms ffs.MinerSelector, dm *dealsModule.Module, wm wallet.Module, ipfs iface.CoreAPI, chain FilChain, l ffs.JobLogger, lsm *lotus.SyncMonitor, minPieceSize uint64, maxParallelDealPreparing int, retrievalNextEventTimeout time.Duration) *FilCold {
5860
return &FilCold{
59-
ms: ms,
60-
dm: dm,
61-
wm: wm,
62-
ipfs: ipfs,
63-
chain: chain,
64-
l: l,
65-
lsm: lsm,
66-
minPieceSize: minPieceSize,
67-
semaphDealPrep: make(chan struct{}, maxParallelDealPreparing),
61+
ms: ms,
62+
dm: dm,
63+
wm: wm,
64+
ipfs: ipfs,
65+
chain: chain,
66+
l: l,
67+
lsm: lsm,
68+
minPieceSize: minPieceSize,
69+
retrNextEventTimeout: retrievalNextEventTimeout,
70+
semaphDealPrep: make(chan struct{}, maxParallelDealPreparing),
6871
}
6972
}
7073

@@ -76,21 +79,39 @@ func (fc *FilCold) Fetch(ctx context.Context, pyCid cid.Cid, piCid *cid.Cid, wad
7679
return ffs.FetchInfo{}, fmt.Errorf("fetching from deal module: %s", err)
7780
}
7881
fc.l.Log(ctx, "Fetching from %s...", miner)
79-
var fundsSpent uint64
80-
var lastMsg string
81-
for e := range events {
82-
if e.Err != "" {
83-
return ffs.FetchInfo{}, fmt.Errorf("event error in retrieval progress: %s", e.Err)
84-
}
85-
strEvent := retrievalmarket.ClientEvents[e.Event]
86-
strDealStatus := retrievalmarket.DealStatuses[e.Status]
87-
fundsSpent = e.FundsSpent.Uint64()
88-
newMsg := fmt.Sprintf("Received %s, total spent: %sFIL (%s/%s)", humanize.IBytes(e.BytesReceived), util.AttoFilToFil(fundsSpent), strEvent, strDealStatus)
89-
if newMsg != lastMsg {
90-
fc.l.Log(ctx, newMsg)
91-
lastMsg = newMsg
82+
83+
var (
84+
fundsSpent uint64
85+
lastMsg string
86+
lastEvent marketevents.RetrievalEvent
87+
)
88+
Loop:
89+
for {
90+
select {
91+
case <-time.After(fc.retrNextEventTimeout):
92+
return ffs.FetchInfo{}, fmt.Errorf("didn't receive events for %d minutes", int64(fc.retrNextEventTimeout.Minutes()))
93+
case e, ok := <-events:
94+
if !ok {
95+
break Loop
96+
}
97+
if e.Err != "" {
98+
return ffs.FetchInfo{}, fmt.Errorf("event error in retrieval progress: %s", e.Err)
99+
}
100+
strEvent := retrievalmarket.ClientEvents[e.Event]
101+
strDealStatus := retrievalmarket.DealStatuses[e.Status]
102+
fundsSpent = e.FundsSpent.Uint64()
103+
newMsg := fmt.Sprintf("Received %s, total spent: %sFIL (%s/%s)", humanize.IBytes(e.BytesReceived), util.AttoFilToFil(fundsSpent), strEvent, strDealStatus)
104+
if newMsg != lastMsg {
105+
fc.l.Log(ctx, newMsg)
106+
lastMsg = newMsg
107+
}
108+
lastEvent = e
92109
}
93110
}
111+
if lastEvent.Status != retrievalmarket.DealStatusCompleted {
112+
return ffs.FetchInfo{}, fmt.Errorf("retrieval failed with status %s and message %s", retrievalmarket.DealStatuses[lastEvent.Status], lastMsg)
113+
}
114+
94115
return ffs.FetchInfo{RetrievedMiner: miner, FundsSpent: fundsSpent}, nil
95116
}
96117

ffs/integrationtest/manager/manager.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@ func NewCustomFFSManager(t require.TestingT, ds datastore.TxnDatastore, cb lotus
8686
l := joblogger.New(txndstr.Wrap(ds, "ffs/joblogger"))
8787
lsm, err := lotus.NewSyncMonitor(cb)
8888
require.NoError(t, err)
89-
cl := filcold.New(ms, dm, nil, ipfsClient, fchain, l, lsm, minimumPieceSize, 1)
89+
cl := filcold.New(ms, dm, nil, ipfsClient, fchain, l, lsm, minimumPieceSize, 1, time.Hour)
9090
hl, err := coreipfs.New(ds, ipfsClient, l)
9191
require.NoError(t, err)
9292
sched, err := scheduler.New(txndstr.Wrap(ds, "ffs/scheduler"), l, hl, cl, 10, time.Minute*10, nil, scheduler.GCConfig{AutoGCInterval: 0})

ffs/interfaces.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ type HotStorage interface {
4242
// Stage adds io.Reader and stage-pins it.
4343
Stage(context.Context, APIID, io.Reader) (cid.Cid, error)
4444

45-
// StageCid stage-pins a cid.
45+
// StageCid pulls Cid data and stage-pin it.
4646
StageCid(context.Context, APIID, cid.Cid) error
4747

4848
// Unpin unpins a Cid.

ffs/manager/manager.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ var (
3333
Hot: ffs.HotConfig{
3434
Enabled: false,
3535
Ipfs: ffs.IpfsConfig{
36-
AddTimeout: 480, // 8 min
36+
AddTimeout: 15 * 60, // 15min
3737
},
3838
},
3939
Cold: ffs.ColdConfig{

0 commit comments

Comments
 (0)