diff --git a/clients/baseclient/client_error.go b/clients/baseclient/client_error.go index d32e5a1..9b2ff72 100644 --- a/clients/baseclient/client_error.go +++ b/clients/baseclient/client_error.go @@ -49,9 +49,10 @@ 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 { @@ -59,6 +60,7 @@ func (e *ErrorResponse) Error() string { } 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 { diff --git a/clients/baseclient/client_util.go b/clients/baseclient/client_util.go index ab146f1..d441575 100644 --- a/clients/baseclient/client_util.go +++ b/clients/baseclient/client_util.go @@ -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) } @@ -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 diff --git a/clients/baseclient/client_util_test.go b/clients/baseclient/client_util_test.go index d4b0321..519378e 100755 --- a/clients/baseclient/client_util_test.go +++ b/clients/baseclient/client_util_test.go @@ -1,6 +1,7 @@ package baseclient import ( + "errors" "fmt" "time" @@ -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)) + }) + }) }) }) @@ -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)) + }) + }) }) }) diff --git a/clients/mtaclient/mta_rest_client.go b/clients/mtaclient/mta_rest_client.go index 606be47..c0422dd 100644 --- a/clients/mtaclient/mta_rest_client.go +++ b/clients/mtaclient/mta_rest_client.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "encoding/json" + "errors" "fmt" "io" "mime/multipart" @@ -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 @@ -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 @@ -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) { diff --git a/clients/mtaclient/mta_rest_client_test.go b/clients/mtaclient/mta_rest_client_test.go new file mode 100644 index 0000000..65a76d5 --- /dev/null +++ b/clients/mtaclient/mta_rest_client_test.go @@ -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{}} +} diff --git a/clients/mtaclient/mtaclient_suite_test.go b/clients/mtaclient/mtaclient_suite_test.go new file mode 100644 index 0000000..0ed2874 --- /dev/null +++ b/clients/mtaclient/mtaclient_suite_test.go @@ -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") +} diff --git a/commands/deploy_command.go b/commands/deploy_command.go index 0baf06b..5ab5837 100644 --- a/commands/deploy_command.go +++ b/commands/deploy_command.go @@ -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 } @@ -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: "", @@ -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." +}