Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ require (
github.com/go-logr/stdr v1.2.2 // indirect
github.com/go-ole/go-ole v1.3.0 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/gorilla/websocket v1.4.2 // indirect
github.com/gorilla/websocket v1.5.3 // indirect
github.com/holiman/uint256 v1.3.2 // indirect
github.com/ipfs/go-block-format v0.2.0 // indirect
github.com/ipfs/go-ipfs-util v0.0.3 // indirect
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -109,8 +109,8 @@ github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+
github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1/go.mod h1:wJfORRmW1u3UXTncJ5qlYoELFm8eSnnEO6hX4iZ3EWY=
github.com/gopherjs/gopherjs v1.17.2 h1:fQnZVsXk8uxXIStYb0N4bGk7jeyTalG/wsZjQ25dO0g=
github.com/gopherjs/gopherjs v1.17.2/go.mod h1:pRRIvn/QzFLrKfvEz3qUuEhtE/zLCWfreZ6J5gM2i+k=
github.com/gorilla/websocket v1.4.2 h1:+/TMaTYc4QFitKJxsQ7Yye35DkWvkdLcvGKqM+x0Ufc=
github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg=
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
github.com/grafana/pyroscope-go v1.2.7 h1:VWBBlqxjyR0Cwk2W6UrE8CdcdD80GOFNutj0Kb1T8ac=
github.com/grafana/pyroscope-go v1.2.7/go.mod h1:o/bpSLiJYYP6HQtvcoVKiE9s5RiNgjYTj1DhiddP2Pc=
github.com/grafana/pyroscope-go/godeltaprof v0.1.9 h1:c1Us8i6eSmkW+Ez05d3co8kasnuOY813tbMN8i/a3Og=
Expand Down
54 changes: 33 additions & 21 deletions pdp/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"net"
"net/http"
"net/url"
"slices"
"strings"
"syscall"
"time"
Expand Down Expand Up @@ -148,6 +149,7 @@ func (c *Client) doWithClient(client *http.Client, req *http.Request, expectStat
}
resp, err := client.Do(req)
if err != nil {
err = redactRequestError(err)
return nil, nil, fmt.Errorf("pdp: %s %s: %w", req.Method, req.URL.Path, err)
}
defer func() { _ = resp.Body.Close() }()
Expand All @@ -170,18 +172,16 @@ func (c *Client) doWithClient(client *http.Client, req *http.Request, expectStat
}
return resp, body, nil
}
for _, s := range expectStatuses {
if resp.StatusCode == s {
return resp, body, nil
}
if slices.Contains(expectStatuses, resp.StatusCode) {
return resp, body, nil
}
return resp, body, newHTTPError(req, resp, body)
}

// isRetryable reports whether the error warrants a retry attempt.
//
// Non-retryable (permanent):
// - context.Canceled / context.DeadlineExceeded
// - caller context cancellation or deadline
// - HTTP 4xx except 429
// - HTTP 501 Not Implemented
// - TLS alert errors (bad cert, expired cert, protocol violations)
Expand All @@ -196,11 +196,11 @@ func (c *Client) doWithClient(client *http.Client, req *http.Request, expectStat
// Unknown error types are NOT retried. Older releases retried optimistically,
// but that masked permanent misconfigurations (bad URL, invalid signer).
// Callers that need broader retries should do so at the business layer.
func isRetryable(err error) bool {
func isRetryable(ctx context.Context, err error) bool {
if err == nil {
return false
}
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
if ctx.Err() != nil || errors.Is(err, context.Canceled) {
return false
}
if httpErr, ok := errors.AsType[*HTTPError](err); ok {
Expand All @@ -215,6 +215,14 @@ func isRetryable(err error) bool {
if _, ok := errors.AsType[tls.AlertError](err); ok {
return false
}
// http.Client and transport timeouts surface as url.Error. The caller's
// context was checked above, so these are safe to retry for idempotent calls.
if urlErr, ok := errors.AsType[*url.Error](err); ok && urlErr.Timeout() {
return true
}
if errors.Is(err, context.DeadlineExceeded) {
return false
}
// Connection-level failures are transient.
if errors.Is(err, syscall.ECONNREFUSED) ||
errors.Is(err, syscall.ECONNRESET) ||
Expand All @@ -239,14 +247,20 @@ func isRetryable(err error) bool {
if dnsErr, ok := errors.AsType[*net.DNSError](err); ok {
return dnsErr.IsTemporary || dnsErr.IsTimeout
}
// url.Error surfaces timeouts (request timeout, idle timeout).
if urlErr, ok := errors.AsType[*url.Error](err); ok {
return urlErr.Timeout()
}
// Unknown error type: do not retry. Safer than optimistic retry.
return false
}

func redactRequestError(err error) error {
urlErr, ok := errors.AsType[*url.Error](err)
if !ok {
return err
}
redacted := *urlErr
redacted.URL = redactURLString(urlErr.URL)
return &redacted
}

// httpRetryDelay returns the delay to wait before the next retry attempt.
// For responses with a Retry-After header (429, 503) the server's value takes
// precedence, capped at maxRetryDelay. Otherwise an exponential backoff
Expand All @@ -263,10 +277,7 @@ func httpRetryDelay(err error, attempt int) time.Duration {
if attempt > maxShift {
attempt = maxShift
}
d := time.Duration(1<<uint(attempt)) * time.Second
if d > maxRetryDelay {
d = maxRetryDelay
}
d := min(time.Duration(1<<uint(attempt))*time.Second, maxRetryDelay)
return d
}

Expand All @@ -280,10 +291,11 @@ func httpRetryDelay(err error, attempt int) time.Duration {
// mutate server state must not be retried here — see postJSON/deleteJSON.
// Long-running and streaming calls should use c.do directly.
func (c *Client) doRetryable(ctx context.Context, makeReq func() (*http.Request, error), expectStatuses ...int) (*http.Response, []byte, error) {
maxRetries := c.maxRetries
if maxRetries < 0 {
maxRetries = 0
}
return c.doRetryableWithClient(ctx, c.httpClient, makeReq, expectStatuses...)
}

func (c *Client) doRetryableWithClient(ctx context.Context, client *http.Client, makeReq func() (*http.Request, error), expectStatuses ...int) (*http.Response, []byte, error) {
maxRetries := max(c.maxRetries, 0)
for attempt := 0; attempt <= maxRetries; attempt++ {
if err := ctx.Err(); err != nil {
return nil, nil, err
Expand All @@ -292,11 +304,11 @@ func (c *Client) doRetryable(ctx context.Context, makeReq func() (*http.Request,
if err != nil {
return nil, nil, err
}
resp, body, err := c.do(req, expectStatuses...)
resp, body, err := c.doWithClient(client, req, expectStatuses...)
if err == nil {
return resp, body, nil
}
if !isRetryable(err) || attempt == maxRetries {
if !isRetryable(ctx, err) || attempt == maxRetries {
return resp, body, err
}
if c.logger != nil {
Expand Down
94 changes: 62 additions & 32 deletions pdp/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -551,10 +551,10 @@ func TestWaitForDataSetCreated(t *testing.T) {
calls++
w.Header().Set("Content-Type", "application/json")
if calls == 1 {
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"pending","dataSetCreated":false,"ok":null}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"pending","dataSetCreated":false,"ok":null}`)
return
}
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"confirmed","dataSetCreated":true,"ok":true,"dataSetId":42}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"confirmed","dataSetCreated":true,"ok":true,"dataSetId":42}`)
}))
status, err := c.WaitForDataSetCreated(context.Background(), c.BaseURL().String()+"pdp/data-sets/created/0x1", 10*time.Millisecond)
if err != nil {
Expand All @@ -569,7 +569,7 @@ func TestGetDataSetCreationStatus_Accepts202(t *testing.T) {
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusAccepted)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"pending","dataSetCreated":false,"ok":null}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"pending","dataSetCreated":false,"ok":null}`)
}))
status, err := c.GetDataSetCreationStatus(context.Background(), c.BaseURL().String()+"pdp/data-sets/created/0x1")
if err != nil {
Expand All @@ -580,33 +580,30 @@ func TestGetDataSetCreationStatus_Accepts202(t *testing.T) {
}
}

func TestWaitForDataSetCreated_ConfirmedFalseStillPending(t *testing.T) {
func TestWaitForDataSetCreated_ConfirmedWithoutResultStillPending(t *testing.T) {
var calls int
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
calls++
w.Header().Set("Content-Type", "application/json")
if calls == 1 {
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"confirmed","dataSetCreated":false,"ok":null}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"confirmed","dataSetCreated":false,"ok":null}`)
return
}
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"confirmed","dataSetCreated":true,"ok":true,"dataSetId":42}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"confirmed","dataSetCreated":true,"ok":true,"dataSetId":42}`)
}))
status, err := c.WaitForDataSetCreated(context.Background(), c.BaseURL().String()+"pdp/data-sets/created/0x1", 10*time.Millisecond)
status, err := c.WaitForDataSetCreated(context.Background(), c.BaseURL().String()+"pdp/data-sets/created/0x1", time.Millisecond)
if err != nil {
t.Fatal(err)
}
if calls < 2 {
t.Fatalf("expected multiple polls, got %d", calls)
}
if status.DataSetID == nil || !status.DataSetID.Equal(types.NewBigInt(42)) {
t.Fatalf("id=%v", status.DataSetID)
if calls != 2 || status.DataSetID == nil || !status.DataSetID.Equal(types.NewBigInt(42)) {
t.Fatalf("calls=%d status=%+v", calls, status)
}
}

func TestWaitForDataSetCreated_Rejected(t *testing.T) {
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x1","service":"svc","txStatus":"rejected","dataSetCreated":false,"ok":false}`)
_, _ = fmt.Fprint(w, `{"createMessageHash":"0x0000000000000000000000000000000000000000000000000000000000000001","service":"svc","txStatus":"rejected","dataSetCreated":false,"ok":false}`)
}))
_, err := c.WaitForDataSetCreated(context.Background(), c.BaseURL().String()+"pdp/data-sets/created/0x1", 10*time.Millisecond)
if !errors.Is(err, ErrTxRejected) {
Expand Down Expand Up @@ -687,13 +684,13 @@ func TestAddPieces(t *testing.T) {
}

func TestAddPieces_MaxBatchSizeAccepted(t *testing.T) {
pcInfo, err := piece.CalculateFromBytes([]byte("hi"))
if err != nil {
t.Fatalf("CalculateFromBytes: %v", err)
}
pieces := make([]AddPieceInput, MaxAddPiecesBatchSize)
for i := range pieces {
pieces[i] = AddPieceInput{PieceCID: pcInfo.CIDv1}
info, err := piece.CalculateFromBytes([]byte{byte(i), 0xa5})
if err != nil {
t.Fatalf("CalculateFromBytes(%d): %v", i, err)
}
pieces[i] = AddPieceInput{PieceCID: info.CIDv1}
}
var gotPieces int
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
Expand All @@ -713,6 +710,42 @@ func TestAddPieces_MaxBatchSizeAccepted(t *testing.T) {
}
}

func TestAddPiecesRejectsDuplicateCanonicalCIDBeforeRequest(t *testing.T) {
info := testPieceInfoV2(t)
requests := 0
c, _ := newTestClient(t, http.HandlerFunc(func(http.ResponseWriter, *http.Request) {
requests++
}))
_, err := c.AddPieces(context.Background(), types.NewBigInt(5), []AddPieceInput{
{PieceCID: info.CIDv1},
{PieceCID: info.CIDv2},
}, []byte{1})
if err == nil || !strings.Contains(err.Error(), "duplicate pieceCID") {
t.Fatalf("error=%v want duplicate pieceCID", err)
}
if requests != 0 {
t.Fatalf("requests=%d want 0", requests)
}
}

func TestAddPiecesAllowsSameCIDAcrossRequests(t *testing.T) {
info := testPieceInfoV2(t)
requests := 0
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
requests++
w.Header().Set("Location", "/pdp/data-sets/5/pieces/added/0xdead000000000000000000000000000000000000000000000000000000000000")
w.WriteHeader(http.StatusCreated)
}))
for range 2 {
if _, err := c.AddPieces(context.Background(), types.NewBigInt(5), []AddPieceInput{{PieceCID: info.CIDv2}}, []byte{1}); err != nil {
t.Fatal(err)
}
}
if requests != 2 {
t.Fatalf("requests=%d want 2", requests)
}
}

func TestAddPieces_TooManyPieces(t *testing.T) {
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
t.Fatal("should not reach server")
Expand Down Expand Up @@ -756,10 +789,10 @@ func TestWaitForPiecesAdded(t *testing.T) {
calls++
w.Header().Set("Content-Type", "application/json")
if calls == 1 {
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"pending","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"pending","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
return
}
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10,11]}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10,11]}`)
}))
status, err := c.WaitForPiecesAdded(context.Background(), c.BaseURL().String()+"status", 10*time.Millisecond)
if err != nil {
Expand All @@ -774,7 +807,7 @@ func TestGetAddPiecesStatus_Accepts202(t *testing.T) {
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusAccepted)
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"pending","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"pending","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
}))
status, err := c.GetAddPiecesStatus(context.Background(), c.BaseURL().String()+"status")
if err != nil {
Expand All @@ -785,26 +818,23 @@ func TestGetAddPiecesStatus_Accepts202(t *testing.T) {
}
}

func TestWaitForPiecesAdded_ConfirmedFalseStillPending(t *testing.T) {
func TestWaitForPiecesAdded_ConfirmedWithoutResultStillPending(t *testing.T) {
var calls int
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
calls++
w.Header().Set("Content-Type", "application/json")
if calls == 1 {
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":null,"piecesAdded":false}`)
return
}
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10,11]}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"confirmed","dataSetId":5,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10]}`)
}))
status, err := c.WaitForPiecesAdded(context.Background(), c.BaseURL().String()+"status", 10*time.Millisecond)
status, err := c.WaitForPiecesAdded(context.Background(), c.BaseURL().String()+"status", time.Millisecond)
if err != nil {
t.Fatal(err)
}
if calls < 2 {
t.Fatalf("expected multiple polls, got %d", calls)
}
if len(status.ConfirmedPieceIDs) != 2 {
t.Fatalf("len=%d", len(status.ConfirmedPieceIDs))
if calls != 2 || len(status.ConfirmedPieceIDs) != 1 {
t.Fatalf("calls=%d status=%+v", calls, status)
}
}

Expand All @@ -825,7 +855,7 @@ func TestWaitForPiecesAdded_404ReturnsHTTPError(t *testing.T) {
func TestGetAddPiecesStatus_LargeUint64DataSetID(t *testing.T) {
c, _ := newTestClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_, _ = fmt.Fprint(w, `{"txHash":"0x1","txStatus":"confirmed","dataSetId":9223372036854775808,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10]}`)
_, _ = fmt.Fprint(w, `{"txHash":"0x0000000000000000000000000000000000000000000000000000000000000001","txStatus":"confirmed","dataSetId":9223372036854775808,"pieceCount":1,"addMessageOk":true,"piecesAdded":true,"confirmedPieceIds":[10]}`)
}))
status, err := c.GetAddPiecesStatus(context.Background(), c.BaseURL().String()+"status")
if err != nil {
Expand Down Expand Up @@ -923,7 +953,7 @@ func TestIsRetryable(t *testing.T) {
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
if got := isRetryable(tc.err); got != tc.want {
if got := isRetryable(context.Background(), tc.err); got != tc.want {
t.Errorf("isRetryable(%v) = %v, want %v", tc.err, got, tc.want)
}
})
Expand Down
Loading