148 lines
4.0 KiB
Go
148 lines
4.0 KiB
Go
package openstack
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"errors"
|
||
|
|
"io"
|
||
|
|
"net"
|
||
|
|
"strings"
|
||
|
|
"time"
|
||
|
|
)
|
||
|
|
|
||
|
|
// DefaultListPageSize is the page size used when the caller passes a
|
||
|
|
// non-positive one.
|
||
|
|
const DefaultListPageSize = 200
|
||
|
|
|
||
|
|
// DefaultListPageRetries is how many times a single failed page request is
|
||
|
|
// retried (so up to DefaultListPageRetries+1 attempts in total).
|
||
|
|
const DefaultListPageRetries = 5
|
||
|
|
|
||
|
|
// pageFetcher fetches one page of floating IPs: at most limit entries that
|
||
|
|
// follow the entry with ID marker ("" = from the start), in a stable server
|
||
|
|
// order. A page shorter than limit is the last one.
|
||
|
|
type pageFetcher func(ctx context.Context, marker string, limit int) ([]FloatingIP, error)
|
||
|
|
|
||
|
|
// pageRetry configures per-page retries. The zero value retries nothing.
|
||
|
|
type pageRetry struct {
|
||
|
|
// Retries is the number of retries after the first failed attempt.
|
||
|
|
Retries int
|
||
|
|
// Sleep waits for d or until ctx is done; nil uses a real timer. Tests
|
||
|
|
// inject a fake to avoid real backoff delays.
|
||
|
|
Sleep func(ctx context.Context, d time.Duration) error
|
||
|
|
}
|
||
|
|
|
||
|
|
// retryBackoff is the delay before retry number attempt (1-based):
|
||
|
|
// 1s, 2s, 4s, 8s, 16s, then capped at 30s.
|
||
|
|
func retryBackoff(attempt int) time.Duration {
|
||
|
|
if attempt < 1 {
|
||
|
|
attempt = 1
|
||
|
|
}
|
||
|
|
d := time.Second << uint(attempt-1)
|
||
|
|
if d > 30*time.Second || d <= 0 {
|
||
|
|
d = 30 * time.Second
|
||
|
|
}
|
||
|
|
return d
|
||
|
|
}
|
||
|
|
|
||
|
|
func sleepCtx(ctx context.Context, d time.Duration) error {
|
||
|
|
t := time.NewTimer(d)
|
||
|
|
defer t.Stop()
|
||
|
|
select {
|
||
|
|
case <-ctx.Done():
|
||
|
|
return ctx.Err()
|
||
|
|
case <-t.C:
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// statusCoder is implemented by gophercloud.ErrUnexpectedResponseCode (and
|
||
|
|
// anything wrapping it): it carries the HTTP status of a failed request.
|
||
|
|
type statusCoder interface{ GetStatusCode() int }
|
||
|
|
|
||
|
|
// IsRetryableListError reports whether a failed page request is worth
|
||
|
|
// repeating: transport-level failures (timeouts, resets, EOF /
|
||
|
|
// RemoteDisconnected), HTTP 429 and HTTP 5xx. Any other 4xx, context
|
||
|
|
// cancellation and unknown errors are not retried.
|
||
|
|
func IsRetryableListError(err error) bool {
|
||
|
|
if err == nil {
|
||
|
|
return false
|
||
|
|
}
|
||
|
|
if errors.Is(err, context.Canceled) {
|
||
|
|
return false
|
||
|
|
}
|
||
|
|
var sc statusCoder
|
||
|
|
if errors.As(err, &sc) {
|
||
|
|
code := sc.GetStatusCode()
|
||
|
|
return code == 429 || code >= 500
|
||
|
|
}
|
||
|
|
if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) {
|
||
|
|
return true
|
||
|
|
}
|
||
|
|
// A deadline is retryable only when it is a per-request network timeout
|
||
|
|
// (net.Error), which is checked below; the caller's own ctx deadline is
|
||
|
|
// filtered out by the paginator via ctx.Err().
|
||
|
|
var ne net.Error
|
||
|
|
if errors.As(err, &ne) {
|
||
|
|
return true
|
||
|
|
}
|
||
|
|
msg := strings.ToLower(err.Error())
|
||
|
|
for _, s := range []string{
|
||
|
|
"remotedisconnected", "remote end closed connection",
|
||
|
|
"connection reset", "connection refused", "broken pipe",
|
||
|
|
"unexpected eof", "eof", "timeout", "tls handshake",
|
||
|
|
} {
|
||
|
|
if strings.Contains(msg, s) {
|
||
|
|
return true
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return false
|
||
|
|
}
|
||
|
|
|
||
|
|
// paginate walks all pages with marker-based pagination, retrying each failed
|
||
|
|
// page per rp, and hands every non-empty page to onPage in server order. It
|
||
|
|
// returns the number of non-empty pages delivered. An error from onPage, a
|
||
|
|
// non-retryable error, exhausted retries or ctx cancellation stop the walk.
|
||
|
|
func paginate(ctx context.Context, pageSize int, rp pageRetry, fetch pageFetcher, onPage func([]FloatingIP) error) (int, error) {
|
||
|
|
if pageSize <= 0 {
|
||
|
|
pageSize = DefaultListPageSize
|
||
|
|
}
|
||
|
|
sleep := rp.Sleep
|
||
|
|
if sleep == nil {
|
||
|
|
sleep = sleepCtx
|
||
|
|
}
|
||
|
|
marker := ""
|
||
|
|
pages := 0
|
||
|
|
for {
|
||
|
|
if err := ctx.Err(); err != nil {
|
||
|
|
return pages, err
|
||
|
|
}
|
||
|
|
var page []FloatingIP
|
||
|
|
var err error
|
||
|
|
for attempt := 0; ; attempt++ {
|
||
|
|
page, err = fetch(ctx, marker, pageSize)
|
||
|
|
if err == nil {
|
||
|
|
break
|
||
|
|
}
|
||
|
|
if ctx.Err() != nil {
|
||
|
|
return pages, ctx.Err()
|
||
|
|
}
|
||
|
|
if attempt >= rp.Retries || !IsRetryableListError(err) {
|
||
|
|
return pages, err
|
||
|
|
}
|
||
|
|
if serr := sleep(ctx, retryBackoff(attempt+1)); serr != nil {
|
||
|
|
return pages, serr
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if len(page) == 0 {
|
||
|
|
return pages, nil
|
||
|
|
}
|
||
|
|
pages++
|
||
|
|
if err := onPage(page); err != nil {
|
||
|
|
return pages, err
|
||
|
|
}
|
||
|
|
if len(page) < pageSize {
|
||
|
|
return pages, nil
|
||
|
|
}
|
||
|
|
marker = page[len(page)-1].ID
|
||
|
|
}
|
||
|
|
}
|