Skip to content
Open
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
8 changes: 5 additions & 3 deletions clients/baseclient/client_error.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,16 +49,18 @@ func BuildErrorResponse(response runtime.ClientResponse, consumer runtime.Consum

// ErrorResponse handles error cases
type ErrorResponse struct {
Code int
Status string
Payload string
Code int
Status string
Payload string
RetryAfter string
}

func (e *ErrorResponse) Error() string {
return fmt.Sprintf("%s (status %d): %v ", e.Status, e.Code, e.Payload)
}

func (e *ErrorResponse) readResponse(response runtime.ClientResponse, consumer runtime.Consumer, formats strfmt.Registry) error {
e.RetryAfter = response.GetHeader("Retry-After")
buf := new(bytes.Buffer)
_, err := buf.ReadFrom(response.Body())
if err != nil {
Expand Down
11 changes: 5 additions & 6 deletions clients/baseclient/client_util.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,6 @@ func CallWithRetry(callback func() (interface{}, error), maxRetriesCount int, re
if !shouldRetry(err) {
return resp, err
}
retryErr, ok := err.(*RetryAfterError)
if ok {
ui.Warn("Retryable error occurred. Retrying after %s", retryErr.Duration)
time.Sleep(retryErr.Duration)
continue
}
ui.Warn("Error occurred: %s. Retrying after: %s.", err.Error(), retryInterval)
time.Sleep(retryInterval)
}
Expand All @@ -29,6 +23,11 @@ func shouldRetry(err error) bool {
if err == nil {
return false
}
// A rate-limit (HTTP 429) response is not retried: the operation stops
// immediately so the caller can surface the server's Retry-After hint.
if _, ok := err.(*RetryAfterError); ok {
return false
}
ae, ok := err.(*ClientError)
if ok {
httpCode := ae.Code
Expand Down
23 changes: 23 additions & 0 deletions clients/baseclient/client_util_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package baseclient

import (
"errors"
"fmt"
"time"

Expand Down Expand Up @@ -64,6 +65,14 @@ var _ = Describe("ClientUtil", func() {
Expect(result).To(Equal(true))
})
})

Context("when the backend rate-limits the request (RetryAfterError)", func() {
It("Retry operation call shouldn't be made", func() {
err := RetryAfterError{Duration: 5 * time.Second}
result := shouldRetry(&err)
Expect(result).To(Equal(false))
})
})
})
})

Expand Down Expand Up @@ -92,6 +101,20 @@ var _ = Describe("ClientUtil", func() {
Expect(result).To(Equal(result))
})
})
Context("when the callback returns a rate-limit error (RetryAfterError)", func() {
It("Should stop after the first attempt without retrying", func() {
attempts := 0
mockCallback := func() (interface{}, error) {
attempts++
return testStruct{}, &RetryAfterError{Duration: 5 * time.Second}
}
_, err := CallWithRetry(mockCallback, 4, time.Duration(0))
Expect(err).To(HaveOccurred())
var retryErr *RetryAfterError
Expect(errors.As(err, &retryErr)).To(BeTrue())
Expect(attempts).To(Equal(1))
})
})
})
})

Expand Down
25 changes: 21 additions & 4 deletions clients/mtaclient/mta_rest_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"mime/multipart"
Expand All @@ -28,6 +29,8 @@ const restBaseURL string = "api/v1/"

const couldNotGetAsyncJobError = "could not get async file upload job"

const defaultRetryAfter = 3 * time.Second

type MtaRestClient struct {
baseclient.BaseClient
client *MtaClient
Expand Down Expand Up @@ -192,6 +195,9 @@ func (c MtaRestClient) StartMtaOperation(operation models.Operation) (ResponseHe
}
resp, err := c.client.Operations.StartMtaOperation(params, token)
if err != nil {
if retryErr := convertToRetryAfterError(err); retryErr != nil {
return ResponseHeader{}, retryErr
}
return ResponseHeader{}, baseclient.NewClientError(err)
}
return ResponseHeader{Location: resp.Location}, nil
Expand Down Expand Up @@ -352,15 +358,26 @@ func (c MtaRestClient) StartUploadMtaArchiveFromUrl(fileUrl string, namespace *s
}

func (c MtaRestClient) handle429(headers http.Header) error {
retryAfter := headers.Get("Retry-After")
return &baseclient.RetryAfterError{Duration: parseRetryAfter(headers.Get("Retry-After"))}
}

func convertToRetryAfterError(err error) *baseclient.RetryAfterError {
var errResp *baseclient.ErrorResponse
if errors.As(err, &errResp) && errResp.Code == http.StatusTooManyRequests {
return &baseclient.RetryAfterError{Duration: parseRetryAfter(errResp.RetryAfter)}
}
return nil
}

func parseRetryAfter(retryAfter string) time.Duration {
if len(retryAfter) == 0 {
retryAfter = "3"
return defaultRetryAfter
}
dur, err := time.ParseDuration(retryAfter + "s")
if err != nil {
return &baseclient.RetryAfterError{Duration: 3 * time.Second}
return defaultRetryAfter
}
return &baseclient.RetryAfterError{Duration: dur}
return dur
}

func (c MtaRestClient) GetAsyncUploadJob(jobId string, namespace *string) (AsyncUploadJobResult, error) {
Expand Down
99 changes: 99 additions & 0 deletions clients/mtaclient/mta_rest_client_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
package mtaclient_test

import (
"bytes"
"errors"
"io"
"net/http"
"time"

"github.com/cloudfoundry-incubator/multiapps-cli-plugin/clients/baseclient"
"github.com/cloudfoundry-incubator/multiapps-cli-plugin/clients/csrf"
"github.com/cloudfoundry-incubator/multiapps-cli-plugin/clients/models"
"github.com/cloudfoundry-incubator/multiapps-cli-plugin/clients/mtaclient"
"github.com/cloudfoundry-incubator/multiapps-cli-plugin/testutil"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
)

var _ = Describe("MtaRestClient", func() {
Describe("StartMtaOperation", func() {
Context("when the backend rate-limits the request with HTTP 429", func() {
It("should return a RetryAfterError carrying the wait time so the caller can stop and report it", func() {
client := newMtaClient(http.StatusTooManyRequests)

_, err := client.StartMtaOperation(models.Operation{})

Expect(err).Should(HaveOccurred())
var retryErr *baseclient.RetryAfterError
Expect(errors.As(err, &retryErr)).Should(BeTrue())
})

It("should carry the server's Retry-After header value", func() {
client := newMtaClientWithHeaders(http.StatusTooManyRequests, http.Header{"Retry-After": []string{"7"}})

_, err := client.StartMtaOperation(models.Operation{})

Expect(err).Should(HaveOccurred())
var retryErr *baseclient.RetryAfterError
Expect(errors.As(err, &retryErr)).Should(BeTrue())
Expect(retryErr.Duration).Should(Equal(7 * time.Second))
})

It("should fall back to the default when Retry-After is absent", func() {
client := newMtaClientWithHeaders(http.StatusTooManyRequests, http.Header{})

_, err := client.StartMtaOperation(models.Operation{})

Expect(err).Should(HaveOccurred())
var retryErr *baseclient.RetryAfterError
Expect(errors.As(err, &retryErr)).Should(BeTrue())
Expect(retryErr.Duration).Should(Equal(3 * time.Second))
})
})

Context("when the backend returns a non-429 error", func() {
It("should return a ClientError, not a RetryAfterError", func() {
client := newMtaClient(http.StatusInternalServerError)

_, err := client.StartMtaOperation(models.Operation{})

Expect(err).Should(HaveOccurred())
var retryErr *baseclient.RetryAfterError
Expect(errors.As(err, &retryErr)).Should(BeFalse())
var clientErr *baseclient.ClientError
Expect(errors.As(err, &clientErr)).Should(BeTrue())
})
})
})
})

func newMtaClient(statusCode int) mtaclient.MtaClientOperations {
tokenFactory := testutil.NewCustomTokenFactory("test-token")
roundTripper := testutil.NewCustomTransport(statusCode)
return mtaclient.NewMtaClient("http://localhost:1000", "test-space-guid", roundTripper, tokenFactory)
}

func newMtaClientWithHeaders(statusCode int, headers http.Header) mtaclient.MtaClientOperations {
tokenFactory := testutil.NewCustomTokenFactory("test-token")
roundTripper := newHeaderTransport(statusCode, headers)
return mtaclient.NewMtaClient("http://localhost:1000", "test-space-guid", roundTripper, tokenFactory)
}

type headerRoundTripperFunc func(*http.Request) (*http.Response, error)

func (fn headerRoundTripperFunc) RoundTrip(req *http.Request) (*http.Response, error) {
return fn(req)
}

func newHeaderTransport(statusCode int, headers http.Header) *csrf.Transport {
transport := headerRoundTripperFunc(func(req *http.Request) (*http.Response, error) {
return &http.Response{
StatusCode: statusCode,
Header: headers,
Body: io.NopCloser(bytes.NewBuffer(nil)),
}, nil
})
userAgentTransport := baseclient.NewUserAgentTransport(transport)
return &csrf.Transport{Delegate: userAgentTransport, Csrf: &csrf.CsrfTokenHelper{}}
}
13 changes: 13 additions & 0 deletions clients/mtaclient/mtaclient_suite_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
package mtaclient_test

import (
"testing"

. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
)

func TestMtaClient(t *testing.T) {
RegisterFailHandler(Fail)
RunSpecs(t, "MtaClient Suite")
}
24 changes: 23 additions & 1 deletion commands/deploy_command.go
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,11 @@ func (c *DeployCommand) executeInternal(positionalArgs []string, dsHost string,
// Create the new process
responseHeader, err := mtaClient.StartMtaOperation(*operation)
if err != nil {
var retryErr *baseclient.RetryAfterError
if errors.As(err, &retryErr) {
ui.Failed(rateLimitMessage(retryErr.Duration))
return Failure
}
ui.Failed("Could not create operation: %s", baseclient.NewClientError(err))
return Failure
}
Expand Down Expand Up @@ -509,7 +514,12 @@ func (c *DeployCommand) uploadFromUrl(url string, mtaClient mtaclient.MtaClientO
func (c *DeployCommand) doUploadFromUrl(encodedFileUrl string, mtaClient mtaclient.MtaClientOperations, namespace string, progressBar *pb.ProgressBar) UploadFromUrlStatus {
responseHeaders, err := mtaClient.StartUploadMtaArchiveFromUrl(encodedFileUrl, &namespace)
if err != nil {
ui.Failed("Could not upload from url: %s", err)
var retryErr *baseclient.RetryAfterError
if errors.As(err, &retryErr) {
ui.Failed(rateLimitMessage(retryErr.Duration))
} else {
ui.Failed("Could not upload from url: %s", err)
}
return UploadFromUrlStatus{
FileId: "",
MtaId: "",
Expand Down Expand Up @@ -832,3 +842,15 @@ func ValidateBooleanFlag(flagName string, flags *flag.FlagSet) error {

return nil
}

// rateLimitMessage builds the terminal message shown when the deploy service rejects a
// request with HTTP 429. The backend sends Retry-After: 0 for active-operation-cap
// rejections (a slot frees when another operation finishes, not on a timer), so only
// include a concrete wait time when the server gave a positive one.
func rateLimitMessage(retryAfter time.Duration) string {
base := "The deploy service is rate-limiting operations for your user or space. Please try again"
if retryAfter > 0 {
return fmt.Sprintf("%s in %s.", base, retryAfter)
}
return base + " later."
}
Loading