Skip to content

Commit c99b724

Browse files
authored
fix(storage): harden provider health checks (#1)
* fix(storage): harden provider health checks * fix(storage): stop probes beyond selection frontier
1 parent ba650c9 commit c99b724

10 files changed

Lines changed: 612 additions & 143 deletions

File tree

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.1.0.
10-
TSSDKRef = "5c7ad6ccc0500d937435a2e0dd7d09549aeccbeb"
9+
// TSSDKRef is the pinned TypeScript SDK commit for synapse-sdk-v1.1.1.
10+
TSSDKRef = "a1d44296ad27b4a2631cb744de95a6a94c8097a7"
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: 4 additions & 4 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.1.0" {
52-
t.Fatalf("synapse-sdk package version = %q, want 1.1.0", sdkPackage.Version)
51+
if sdkPackage.Version != "1.1.1" {
52+
t.Fatalf("synapse-sdk package version = %q, want 1.1.1", sdkPackage.Version)
5353
}
5454
coreDep := sdkPackage.Dependencies["@filoz/synapse-core"]
5555
if coreDep == "" {
@@ -63,8 +63,8 @@ func TestLocalTSSDKBaseline(t *testing.T) {
6363
if corePackage.Name != "@filoz/synapse-core" {
6464
t.Fatalf("synapse-core package name = %q, want @filoz/synapse-core", corePackage.Name)
6565
}
66-
if corePackage.Version != "0.7.0" {
67-
t.Fatalf("synapse-core package version = %q, want 0.7.0", corePackage.Version)
66+
if corePackage.Version != "0.7.1" {
67+
t.Fatalf("synapse-core package version = %q, want 0.7.1", corePackage.Version)
6868
}
6969
t.Logf("local TS baseline packages: %s@%s depends on %s@%s via %s", sdkPackage.Name, sdkPackage.Version, corePackage.Name, corePackage.Version, coreDep)
7070

payments/coverage_test.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -245,6 +245,29 @@ func TestGetSettlementAmounts_DecodesTuple(t *testing.T) {
245245
}
246246
}
247247

248+
func TestReadsDoNotUseSignerAddress(t *testing.T) {
249+
s, mb := newTestService(t)
250+
if s.Account() == (common.Address{}) {
251+
t.Fatal("test service has no signer account")
252+
}
253+
mb.rejectNonZeroCallFrom = true
254+
mb.setFilPayReply(t, filPayAddr, "accounts",
255+
big.NewInt(1000), big.NewInt(200), big.NewInt(5), big.NewInt(42))
256+
mb.setFilPayReply(t, filPayAddr, "getAccountInfoIfSettled",
257+
big.NewInt(0), big.NewInt(1000), big.NewInt(123), big.NewInt(5))
258+
mb.setFilPayReply(t, filPayAddr, "settleRail",
259+
big.NewInt(100), big.NewInt(90), big.NewInt(5), big.NewInt(5),
260+
big.NewInt(50), "ok",
261+
)
262+
263+
if _, err := s.AccountInfo(context.Background(), tokenAddr, s.Account()); err != nil {
264+
t.Fatalf("AccountInfo: %v", err)
265+
}
266+
if _, err := s.GetSettlementAmounts(context.Background(), sdktypes.NewBigInt(3), big.NewInt(50)); err != nil {
267+
t.Fatalf("GetSettlementAmounts: %v", err)
268+
}
269+
}
270+
248271
func TestGetSettlementAmounts_InvalidRailID(t *testing.T) {
249272
s, _ := newTestService(t)
250273
if _, err := s.GetSettlementAmounts(context.Background(), sdktypes.NewBigInt(0), nil); !errors.Is(err, ErrInvalidArgument) {

payments/service_test.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,8 @@ type mockBackend struct {
3636
// errors keyed the same way
3737
errs map[string]error
3838
// callReplyFn, when set, can override a contract call reply.
39-
callReplyFn func(contractHex, method string, data []byte) ([]byte, bool, error)
39+
callReplyFn func(contractHex, method string, data []byte) ([]byte, bool, error)
40+
rejectNonZeroCallFrom bool
4041

4142
// per-account balances; nil -> 0
4243
balances map[common.Address]*big.Int
@@ -91,6 +92,9 @@ func (m *mockBackend) CodeAt(_ context.Context, _ common.Address, _ *big.Int) ([
9192
func (m *mockBackend) CallContract(_ context.Context, call ethereum.CallMsg, blockNumber *big.Int) ([]byte, error) {
9293
m.mu.Lock()
9394
defer m.mu.Unlock()
95+
if m.rejectNonZeroCallFrom && call.From != (common.Address{}) {
96+
return nil, errors.New("non-zero eth_call sender")
97+
}
9498
if len(call.Data) < 4 || call.To == nil {
9599
return nil, errors.New("calldata too short")
96100
}

services.go

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -126,16 +126,10 @@ func (c *Client) initServices() error {
126126
}
127127

128128
resolver, err := storage.NewServiceResolver(storage.ServiceResolverOptions{
129-
Payer: c.evmSigner.EVMAddress(),
130-
SPRegistry: spReg,
131-
WarmStorage: ws,
132-
ProviderPing: func(ctx context.Context, serviceURL string) error {
133-
pdpClient, err := c.newPDPClient(serviceURL, pdp.WithMaxRetries(0))
134-
if err != nil {
135-
return err
136-
}
137-
return pdpClient.Ping(ctx)
138-
},
129+
Payer: c.evmSigner.EVMAddress(),
130+
SPRegistry: spReg,
131+
WarmStorage: ws,
132+
ProviderPing: c.pingProvider,
139133
NewContext: func(sel storage.ResolvedUploadContext, opts *storage.UploadOptions) (*storage.Context, error) {
140134
pdpClient, err := c.newPDPClient(sel.Provider.ServiceURL)
141135
if err != nil {
@@ -224,6 +218,14 @@ func (c *Client) newPDPClient(serviceURL string, opts ...pdp.Option) (*pdp.Clien
224218
return pdp.New(serviceURL, pdpOpts...)
225219
}
226220

221+
func (c *Client) pingProvider(ctx context.Context, serviceURL string) error {
222+
pdpClient, err := c.newPDPClient(serviceURL, pdp.WithMaxRetries(2))
223+
if err != nil {
224+
return err
225+
}
226+
return pdpClient.Ping(ctx)
227+
}
228+
227229
// WarmStorage returns the [warmstorage.Service].
228230
func (c *Client) WarmStorage() *warmstorage.Service {
229231
return c.warmStorage

storage/selector.go

Lines changed: 139 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,9 @@ const selectorListPageSize = 100
2727

2828
const selectorRetryInitialDelay = 200 * time.Millisecond
2929

30-
const providerPingTimeout = 2 * time.Second
30+
const providerPingTimeout = 8 * time.Second
31+
32+
const providerPingConcurrency = 16
3133

3234
const resolveConcurrency = 10
3335

@@ -43,6 +45,19 @@ type autoSelectCandidate struct {
4345
metadata map[string]string
4446
}
4547

48+
type providerProbeState uint8
49+
50+
const (
51+
providerProbePending providerProbeState = iota
52+
providerProbeHealthy
53+
providerProbeUnhealthy
54+
)
55+
56+
type providerProbeResult struct {
57+
index int
58+
err error
59+
}
60+
4661
// PDPProviderSource is the subset of spregistry.Service used by ServiceResolver.
4762
type PDPProviderSource interface {
4863
GetPDPProvider(context.Context, types.BigInt) (*spregistry.PDPProvider, error)
@@ -84,9 +99,12 @@ type ServiceResolverOptions struct {
8499
DataSetValidator DataSetValidator
85100
DataSetDetails DataSetDetailsCatalog
86101
// ProviderPing checks one automatically selected provider endpoint. The
87-
// resolver passes a context with a two-second deadline; the implementation
88-
// is responsible for honoring it. nil uses an unretried [pdp.Client.Ping]
89-
// request.
102+
// resolver may call it concurrently for up to 16 candidates per selection;
103+
// concurrent selections have independent limits. Each call receives a
104+
// context whose deadline is no later than eight seconds; a shorter parent
105+
// deadline takes precedence. Implementations must be safe for concurrent use
106+
// and honor the context. nil uses [pdp.Client.Ping] with up to two retries for
107+
// transient failures.
90108
ProviderPing func(context.Context, string) error
91109
NewContext ContextFactory // called per-provider to construct a managed Context
92110
}
@@ -122,8 +140,8 @@ var (
122140
// NewServiceResolver constructs a ServiceResolver. Payer, SPRegistry,
123141
// WarmStorage, and NewContext are required. DataSetValidator, DataSetDetails,
124142
// and active-piece reads are optional and auto-detected from WarmStorage when
125-
// available and configured. ProviderPing is optional and defaults to an
126-
// unretried PDP ping request.
143+
// available and configured. ProviderPing is optional and defaults to a PDP
144+
// ping with up to two retries for transient failures.
127145
func NewServiceResolver(opts ServiceResolverOptions) (*ServiceResolver, error) {
128146
if opts.Payer == (common.Address{}) {
129147
return nil, fmt.Errorf("storage.NewServiceResolver: %w: zero payer", ErrInvalidArgument)
@@ -169,7 +187,7 @@ func NewServiceResolver(opts ServiceResolverOptions) (*ServiceResolver, error) {
169187
}
170188

171189
func defaultProviderPing(ctx context.Context, serviceURL string) error {
172-
client, err := pdp.New(serviceURL, pdp.WithMaxRetries(0))
190+
client, err := pdp.New(serviceURL, pdp.WithMaxRetries(2))
173191
if err != nil {
174192
return err
175193
}
@@ -441,29 +459,7 @@ func (r *ServiceResolver) autoSelect(ctx context.Context, opts *UploadOptions, e
441459
providersWithDetails = detailedCandidateProviders(detailedDataSets, selectableProviders)
442460
}
443461
requestedMetadata := dataSetMetadataFromOptions(opts)
444-
selected := make([]ResolvedUploadContext, 0, min(count, len(providers)))
445-
failedProviderIDs := make([]types.BigInt, 0, min(count, len(providers)))
446-
probeCandidate := func(candidate autoSelectCandidate) (bool, error) {
447-
pingCtx, cancel := context.WithTimeout(ctx, providerPingTimeout)
448-
pingErr := r.providerPing(pingCtx, candidate.provider.Offering.ServiceURL)
449-
cancel()
450-
if ctxErr := ctx.Err(); ctxErr != nil {
451-
return false, fmt.Errorf("storage.ServiceResolver.ResolveUploadContexts: health-check provider %s: %w", candidate.provider.Info.ID.String(), ctxErr)
452-
}
453-
if pingErr != nil {
454-
failedProviderIDs = append(failedProviderIDs, candidate.provider.Info.ID)
455-
return false, nil
456-
}
457-
selected = append(selected, buildResolvedUploadContext(
458-
*candidate.provider,
459-
candidate.dataSetID,
460-
candidate.clientDataSetID,
461-
candidate.metadata,
462-
))
463-
return len(selected) == count, nil
464-
}
465-
466-
deferWithoutDataSet := len(providersWithDetails) > 0
462+
withDataSet := make([]autoSelectCandidate, 0, min(count, len(providers)))
467463
withoutDataSet := make([]autoSelectCandidate, 0, min(count, len(providers)))
468464
seenProviderIDs := make(map[string]struct{}, min(count, len(providers)))
469465
for i := range providers {
@@ -481,49 +477,117 @@ func (r *ServiceResolver) autoSelect(ctx context.Context, opts *UploadOptions, e
481477
var metadata map[string]string
482478
if providerDataSets, ok := providersWithDetails[providerKey]; ok {
483479
dataSetID, clientDataSetID, metadata = selectMatchingDetailedDataSet(provider.Info.ID, providerDataSets, requestedMetadata)
484-
delete(providersWithDetails, providerKey)
485480
}
486481
candidate := autoSelectCandidate{
487482
provider: provider,
488483
dataSetID: dataSetID,
489484
clientDataSetID: clientDataSetID,
490485
metadata: metadata,
491486
}
492-
switch {
493-
case dataSetID != nil:
494-
done, err := probeCandidate(candidate)
495-
if err != nil {
496-
return nil, err
497-
}
498-
if done {
499-
return selected, nil
500-
}
501-
case deferWithoutDataSet:
487+
if dataSetID != nil {
488+
withDataSet = append(withDataSet, candidate)
489+
} else {
502490
candidate.metadata = requestedMetadata
503491
withoutDataSet = append(withoutDataSet, candidate)
504-
default:
505-
candidate.metadata = requestedMetadata
506-
done, err := probeCandidate(candidate)
507-
if err != nil {
508-
return nil, err
509-
}
510-
if done {
511-
return selected, nil
492+
}
493+
}
494+
candidates := slices.Concat(withDataSet, withoutDataSet)
495+
return r.selectHealthyCandidates(ctx, candidates, count)
496+
}
497+
498+
func (r *ServiceResolver) selectHealthyCandidates(ctx context.Context, candidates []autoSelectCandidate, count int) ([]ResolvedUploadContext, error) {
499+
if len(candidates) == 0 {
500+
return nil, errors.New("storage.ServiceResolver.ResolveUploadContexts: no remaining providers")
501+
}
502+
if err := ctx.Err(); err != nil {
503+
return nil, fmt.Errorf("storage.ServiceResolver.ResolveUploadContexts: health-check providers: %w", err)
504+
}
505+
506+
probeCtx, cancelProbes := context.WithCancel(ctx)
507+
defer cancelProbes()
508+
states := make([]providerProbeState, len(candidates))
509+
results := make(chan providerProbeResult, min(providerPingConcurrency, len(candidates)))
510+
probeCancels := make([]context.CancelFunc, len(candidates))
511+
nextIndex := 0
512+
inFlight := 0
513+
probeLimit := len(candidates)
514+
stopping := false
515+
selectionDetermined := false
516+
517+
startProbes := func() {
518+
for !stopping && inFlight < providerPingConcurrency && nextIndex < probeLimit {
519+
index := nextIndex
520+
nextIndex++
521+
inFlight++
522+
pingCtx, cancel := context.WithTimeout(probeCtx, providerPingTimeout)
523+
probeCancels[index] = cancel
524+
go func() {
525+
err := r.providerPing(pingCtx, candidates[index].provider.Offering.ServiceURL)
526+
cancel()
527+
results <- providerProbeResult{index: index, err: err}
528+
}()
529+
}
530+
}
531+
cancelProbesFrom := func(index int) {
532+
for i := index; i < nextIndex; i++ {
533+
if probeCancels[i] != nil {
534+
probeCancels[i]()
512535
}
513536
}
537+
}
514538

515-
if deferWithoutDataSet && len(providersWithDetails) == 0 {
516-
deferWithoutDataSet = false
517-
for _, deferred := range withoutDataSet {
518-
done, err := probeCandidate(deferred)
519-
if err != nil {
520-
return nil, err
539+
startProbes()
540+
for inFlight > 0 {
541+
result := <-results
542+
inFlight--
543+
probeCancels[result.index] = nil
544+
if !stopping {
545+
if ctxErr := ctx.Err(); ctxErr != nil {
546+
stopping = true
547+
cancelProbes()
548+
} else if result.index < probeLimit {
549+
if result.err == nil {
550+
states[result.index] = providerProbeHealthy
551+
} else {
552+
states[result.index] = providerProbeUnhealthy
553+
}
554+
frontier, ready := providerSelectionProgress(states, count)
555+
if frontier >= 0 && frontier+1 < probeLimit {
556+
// Later candidates cannot enter the stable ranking once this
557+
// prefix already contains the requested healthy copies.
558+
probeLimit = frontier + 1
559+
cancelProbesFrom(probeLimit)
521560
}
522-
if done {
523-
return selected, nil
561+
if ready {
562+
selectionDetermined = true
563+
stopping = true
564+
cancelProbes()
524565
}
525566
}
526-
withoutDataSet = nil
567+
}
568+
startProbes()
569+
}
570+
571+
if err := ctx.Err(); err != nil && !selectionDetermined {
572+
return nil, fmt.Errorf("storage.ServiceResolver.ResolveUploadContexts: health-check providers: %w", err)
573+
}
574+
selected := make([]ResolvedUploadContext, 0, min(count, len(candidates)))
575+
failedProviderIDs := make([]types.BigInt, 0, min(count, len(candidates)))
576+
for i, state := range states {
577+
switch state {
578+
case providerProbeHealthy:
579+
candidate := candidates[i]
580+
selected = append(selected, buildResolvedUploadContext(
581+
*candidate.provider,
582+
candidate.dataSetID,
583+
candidate.clientDataSetID,
584+
candidate.metadata,
585+
))
586+
if len(selected) == count {
587+
return selected, nil
588+
}
589+
case providerProbeUnhealthy:
590+
failedProviderIDs = append(failedProviderIDs, candidates[i].provider.Info.ID)
527591
}
528592
}
529593

@@ -540,6 +604,23 @@ func (r *ServiceResolver) autoSelect(ctx context.Context, opts *UploadOptions, e
540604
return selected, nil
541605
}
542606

607+
func providerSelectionProgress(states []providerProbeState, count int) (frontier int, ready bool) {
608+
healthy := 0
609+
pending := false
610+
for i, state := range states {
611+
if state == providerProbePending {
612+
pending = true
613+
}
614+
if state == providerProbeHealthy {
615+
healthy++
616+
if healthy == count {
617+
return i, !pending
618+
}
619+
}
620+
}
621+
return -1, false
622+
}
623+
543624
func formatProviderIDs(ids []types.BigInt) string {
544625
values := make([]string, len(ids))
545626
for i, id := range ids {

0 commit comments

Comments
 (0)