Skip to content

Commit 0cf2092

Browse files
authored
feat(storage): prefer active data sets for explicit providers
1 parent 1082357 commit 0cf2092

13 files changed

Lines changed: 1087 additions & 62 deletions

File tree

internal/adapters/pdpverifier.go

Lines changed: 4 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,8 @@ import (
44
"context"
55
"fmt"
66
"math/big"
7-
"strings"
87

98
"github.com/ethereum/go-ethereum/accounts/abi/bind"
10-
"github.com/ethereum/go-ethereum/crypto"
119
"github.com/ethereum/go-ethereum/ethclient"
1210
"github.com/ipfs/go-cid"
1311

@@ -57,7 +55,7 @@ func (a *pdpVerifierReader) FindPieceIdsByCid(ctx context.Context, dataSetID sdk
5755
new(big.Int).SetUint64(limit),
5856
)
5957
if err != nil {
60-
if isDataSetUnavailable(err) {
58+
if pdpverifier.IsDataSetUnavailable(err) {
6159
return []sdktypes.BigInt{}, nil
6260
}
6361
return nil, fmt.Errorf("adapters.pdpVerifierReader.FindPieceIdsByCid: %w", err)
@@ -72,7 +70,7 @@ func (a *pdpVerifierReader) FindPieceIdsByCid(ctx context.Context, dataSetID sdk
7270
func (a *pdpVerifierReader) GetScheduledRemovals(ctx context.Context, dataSetID sdktypes.BigInt) ([]sdktypes.BigInt, error) {
7371
raw, err := a.caller.GetScheduledRemovals(&bind.CallOpts{Context: ctx}, dataSetID.Big())
7472
if err != nil {
75-
if isDataSetUnavailable(err) {
73+
if pdpverifier.IsDataSetUnavailable(err) {
7674
return []sdktypes.BigInt{}, nil
7775
}
7876
return nil, fmt.Errorf("adapters.pdpVerifierReader.GetScheduledRemovals: %w", err)
@@ -87,7 +85,7 @@ func (a *pdpVerifierReader) GetScheduledRemovals(ctx context.Context, dataSetID
8785
func (a *pdpVerifierReader) GetNextChallengeEpoch(ctx context.Context, dataSetID sdktypes.BigInt) (*big.Int, error) {
8886
v, err := a.caller.GetNextChallengeEpoch(&bind.CallOpts{Context: ctx}, dataSetID.Big())
8987
if err != nil {
90-
if isDataSetUnavailable(err) {
88+
if pdpverifier.IsDataSetUnavailable(err) {
9189
return nil, nil
9290
}
9391
return nil, fmt.Errorf("adapters.pdpVerifierReader.GetNextChallengeEpoch: %w", err)
@@ -113,7 +111,7 @@ func (a *pdpVerifierReader) BlockNumber(ctx context.Context) (uint64, error) {
113111
func (a *pdpVerifierReader) GetDataSetSizeBytes(ctx context.Context, dataSetID sdktypes.BigInt) (*big.Int, error) {
114112
leafCount, err := a.caller.GetDataSetLeafCount(&bind.CallOpts{Context: ctx}, dataSetID.Big())
115113
if err != nil {
116-
if isDataSetUnavailable(err) {
114+
if pdpverifier.IsDataSetUnavailable(err) {
117115
return new(big.Int), nil
118116
}
119117
return nil, fmt.Errorf("adapters.pdpVerifierReader.GetDataSetSizeBytes: %w", err)
@@ -124,36 +122,6 @@ func (a *pdpVerifierReader) GetDataSetSizeBytes(ctx context.Context, dataSetID s
124122
return new(big.Int).Mul(leafCount, big.NewInt(32)), nil
125123
}
126124

127-
var (
128-
dataSetNotFoundSelector = errorSelector("DataSetNotFound()")
129-
dataSetNotLiveSelector = errorSelector("DataSetNotLive()")
130-
)
131-
132-
// isDataSetUnavailable reports whether err is the PDPVerifier revert raised
133-
// for terminated, missing, or otherwise non-readable data sets.
134-
func isDataSetUnavailable(err error) bool {
135-
if err == nil {
136-
return false
137-
}
138-
msg := err.Error()
139-
if strings.Contains(msg, "Data set not live") ||
140-
strings.Contains(msg, "DataSetNotLive") ||
141-
strings.Contains(msg, "DataSetNotFound") {
142-
return true
143-
}
144-
data, ok := ethclient.RevertErrorData(err)
145-
if !ok || len(data) < 4 {
146-
return false
147-
}
148-
selector := [4]byte{data[0], data[1], data[2], data[3]}
149-
return selector == dataSetNotFoundSelector || selector == dataSetNotLiveSelector
150-
}
151-
152-
func errorSelector(signature string) [4]byte {
153-
hash := crypto.Keccak256([]byte(signature))
154-
return [4]byte{hash[0], hash[1], hash[2], hash[3]}
155-
}
156-
157125
func dedupeBigInts(values []sdktypes.BigInt) []sdktypes.BigInt {
158126
if len(values) == 0 {
159127
return values
Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
1-
// Package pdpverifier contains the generated Go bindings for the
2-
// PDPVerifier smart contract.
1+
// Package pdpverifier contains generated Go bindings and internal error
2+
// helpers for the PDPVerifier smart contract.
33
//
44
// Regenerate with: go generate ./internal/contracts/...
55
package pdpverifier
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
package pdpverifier
2+
3+
import (
4+
"strings"
5+
6+
"github.com/ethereum/go-ethereum/crypto"
7+
"github.com/ethereum/go-ethereum/ethclient"
8+
)
9+
10+
var (
11+
dataSetNotFoundSelector = pdpErrorSelector("DataSetNotFound()")
12+
dataSetNotLiveSelector = pdpErrorSelector("DataSetNotLive()")
13+
)
14+
15+
// IsDataSetUnavailable reports whether err is the PDPVerifier revert raised
16+
// for terminated, missing, or otherwise non-readable data sets.
17+
func IsDataSetUnavailable(err error) bool {
18+
if err == nil {
19+
return false
20+
}
21+
msg := err.Error()
22+
if strings.Contains(msg, "Data set not live") ||
23+
strings.Contains(msg, "DataSetNotLive") ||
24+
strings.Contains(msg, "DataSetNotFound") {
25+
return true
26+
}
27+
data, ok := ethclient.RevertErrorData(err)
28+
if !ok || len(data) < 4 {
29+
return false
30+
}
31+
selector := [4]byte{data[0], data[1], data[2], data[3]}
32+
return selector == dataSetNotFoundSelector || selector == dataSetNotLiveSelector
33+
}
34+
35+
func pdpErrorSelector(signature string) [4]byte {
36+
hash := crypto.Keccak256([]byte(signature))
37+
return [4]byte{hash[0], hash[1], hash[2], hash[3]}
38+
}
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
package pdpverifier
2+
3+
import (
4+
"errors"
5+
"testing"
6+
7+
"github.com/ethereum/go-ethereum/common/hexutil"
8+
"github.com/ethereum/go-ethereum/crypto"
9+
)
10+
11+
type revertDataError struct {
12+
data string
13+
}
14+
15+
func (e revertDataError) Error() string {
16+
return "execution reverted"
17+
}
18+
19+
func (e revertDataError) ErrorCode() int {
20+
return 3
21+
}
22+
23+
func (e revertDataError) ErrorData() interface{} {
24+
return e.data
25+
}
26+
27+
func customError(name string) error {
28+
hash := crypto.Keccak256([]byte(name + "()"))
29+
return revertDataError{data: hexutil.Encode(hash[:4])}
30+
}
31+
32+
func TestIsDataSetUnavailable(t *testing.T) {
33+
t.Parallel()
34+
35+
tests := []struct {
36+
name string
37+
err error
38+
want bool
39+
}{
40+
{name: "nil", err: nil, want: false},
41+
{name: "ordinary error", err: errors.New("connection reset"), want: false},
42+
{name: "legacy not live string", err: errors.New("execution reverted: Data set not live"), want: true},
43+
{name: "custom not live string", err: errors.New("execution reverted: DataSetNotLive()"), want: true},
44+
{name: "custom not found string", err: errors.New("execution reverted: DataSetNotFound()"), want: true},
45+
{name: "custom not live selector", err: customError("DataSetNotLive"), want: true},
46+
{name: "custom not found selector", err: customError("DataSetNotFound"), want: true},
47+
{name: "other selector", err: customError("CleanupDepositRequired"), want: false},
48+
}
49+
50+
for _, tt := range tests {
51+
t.Run(tt.name, func(t *testing.T) {
52+
t.Parallel()
53+
if got := IsDataSetUnavailable(tt.err); got != tt.want {
54+
t.Fatalf("IsDataSetUnavailable() = %t, want %t", got, tt.want)
55+
}
56+
})
57+
}
58+
}

internal/upstream/upstream.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@ const (
66
TSSDKRepo = "FilOzone/synapse-sdk"
77
// TSSDKLocalDir is the local TypeScript SDK checkout path.
88
TSSDKLocalDir = "synapse-sdk"
9-
// TSSDKRef is the pinned TypeScript SDK commit for synapse-sdk-v1.0.1.
10-
TSSDKRef = "e044ddf62210f505d918e723a7f4cd62f4e49ab0"
9+
// TSSDKRef is the pinned TypeScript SDK commit for synapse-sdk-v1.1.0.
10+
TSSDKRef = "5c7ad6ccc0500d937435a2e0dd7d09549aeccbeb"
1111
// FilecoinServicesRepo is the upstream contract ABI repository.
1212
FilecoinServicesRepo = "FilOzone/filecoin-services"
1313
// FilecoinServicesRef is the pinned contract ABI commit.

internal/upstream/upstream_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,8 +48,8 @@ func TestLocalTSSDKBaseline(t *testing.T) {
4848
if sdkPackage.Name != "@filoz/synapse-sdk" {
4949
t.Fatalf("synapse-sdk package name = %q, want @filoz/synapse-sdk", sdkPackage.Name)
5050
}
51-
if sdkPackage.Version != "1.0.1" {
52-
t.Fatalf("synapse-sdk package version = %q, want 1.0.1", sdkPackage.Version)
51+
if sdkPackage.Version != "1.1.0" {
52+
t.Fatalf("synapse-sdk package version = %q, want 1.1.0", sdkPackage.Version)
5353
}
5454
coreDep := sdkPackage.Dependencies["@filoz/synapse-core"]
5555
if coreDep == "" {

storage/doc.go

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -45,11 +45,15 @@
4545
// The Service handles orchestration of the full multi-copy flow, while
4646
// ServiceResolver implements both resolver contracts. It reuses provider-local
4747
// datasets only when metadata matches exactly and the warmstorage-approved
48-
// provider set intersects active PDP providers from spregistry. Automatically
49-
// selected providers must also pass a bounded PDP health check before a context
50-
// is created; explicitly selected provider and data set IDs are not probed. If
51-
// only some automatically selected providers are healthy, the resolver returns
52-
// the healthy subset so upload callers can surface partial-copy results.
48+
// provider set intersects active PDP providers from spregistry. For explicit
49+
// provider IDs, the default warmstorage service checks matching data sets with
50+
// bounded concurrency, preferring the oldest one with active pieces and falling
51+
// back to the oldest empty match. Custom catalogs without active-piece reads
52+
// retain the original oldest-match behavior. Automatically selected providers
53+
// must also pass a bounded PDP health check before a context is created;
54+
// explicitly selected provider and data set IDs are not probed. If only some
55+
// automatically selected providers are healthy, the resolver returns the
56+
// healthy subset so upload callers can surface partial-copy results.
5357
// Existing data sets that cannot accept writes surface typed errors such as
5458
// DataSetPDPPaymentTerminatedError. Use errors.AsType to access fields like
5559
// PDPEndEpoch.

0 commit comments

Comments
 (0)