At 03:30 the store's collector held the registry still while the node-engine's reconcile on that machine was fetching bundle blobs from it: every archive failed 'connection refused' and the machine was held until the next pass. Other machines can meet the same window. A scheduled step now opens its window only once no apply is in flight here (the apply lock is taken just to write the record, so a push still never queues behind the window). An apply whose fetch the store does not answer waits for a window open on its own machine to close, and elsewhere retries with backoff within one bounded budget per apply, well past the window's length; an answer such as 404 still fails at once.
160 lines
6.0 KiB
Go
160 lines
6.0 KiB
Go
package apply
|
|
|
|
// A fetch from the artifact store while the store is away on purpose (novox/hq issue 291, ADR 0189).
|
|
//
|
|
// Every apply fetches each archive and process it declares, to know its bytes by digest. Each night
|
|
// the store's collector holds the store's server still for its maintenance window — seconds today,
|
|
// minutes on a larger store — and an apply that reached its archives in that window failed every one
|
|
// of them with "connection refused" and held the machine until its next pass. Harmless once, and
|
|
// wrong: a planned window must not fail an apply anywhere.
|
|
//
|
|
// **What the apply does about it, in the simplest form that is right on every machine:**
|
|
//
|
|
// - On the store's own machine the window is known exactly: it is the window record every apply
|
|
// there already reads (maintenance_window.go). A fetch the store does not answer while a window
|
|
// is open on this machine waits for the windows to close — polled, bounded by windowWaitAtMost —
|
|
// and is tried again.
|
|
// - On every other machine the store's window is not known, and is not published: a fetch the
|
|
// store does not answer is retried with backoff for at most storeAwayAtMost, **once per apply** —
|
|
// the budget is shared by every fetch of that apply, so a store that is really down costs one
|
|
// bounded wait, not one per archive. The bound is many times the window the collector holds
|
|
// (twelve seconds measured), and whatever still fails is retried by the engine's next pass
|
|
// anyway: the bound only decides whether one pass fails, never whether the machine converges.
|
|
//
|
|
// Only a store that does not answer is waited for: refused, reset, cut short, or a gateway saying
|
|
// the server behind it is away. An answer — a missing blob, the wrong bytes — fails at once, as before.
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"net/http"
|
|
"strings"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
)
|
|
|
|
// How long an apply waits for the windows open on its own machine to close before a fetch the store
|
|
// does not answer fails. Longer than any collection this mesh runs; far shorter than windowAtMost,
|
|
// which is how long a record nothing closed is believed.
|
|
var windowWaitAtMost = 20 * time.Minute
|
|
|
|
// How long, in all, one apply retries a store that does not answer when no window on this machine
|
|
// explains it: on every machine but the store's, the store's own window.
|
|
var storeAwayAtMost = 3 * time.Minute
|
|
|
|
// How often an open window is looked at again, and the first and longest backoff outside one.
|
|
// Variables so a test runs the 03:30 sequence in milliseconds.
|
|
var (
|
|
windowPoll = 2 * time.Second
|
|
storeBackoff = time.Second
|
|
storeBackoffMax = 15 * time.Second
|
|
)
|
|
|
|
// statusError is the store answering with a status other than 200.
|
|
type statusError struct {
|
|
source string
|
|
code int
|
|
status string
|
|
}
|
|
|
|
func (e statusError) Error() string { return fmt.Sprintf("%s answered %s", e.source, e.status) }
|
|
|
|
// storeAwayErr is whether a fetch failed because nothing answered for the store, rather than because
|
|
// the store answered no.
|
|
func storeAwayErr(err error) bool {
|
|
var s statusError
|
|
if errors.As(err, &s) {
|
|
return s.code == http.StatusBadGateway || s.code == http.StatusServiceUnavailable ||
|
|
s.code == http.StatusGatewayTimeout
|
|
}
|
|
if errors.Is(err, syscall.ECONNREFUSED) || errors.Is(err, syscall.ECONNRESET) ||
|
|
errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.EOF) {
|
|
return true
|
|
}
|
|
var n net.Error
|
|
return errors.As(err, &n) && n.Timeout()
|
|
}
|
|
|
|
// storeAway is one apply's patience with a store that does not answer.
|
|
type storeAway struct {
|
|
windows *Windows
|
|
log func(string)
|
|
|
|
mu sync.Mutex
|
|
windowWaited time.Duration
|
|
left time.Duration
|
|
said map[string]bool
|
|
}
|
|
|
|
// newStoreAway is the patience of one apply on a machine whose windows are recorded in windows (nil:
|
|
// none known, as on every machine but the store's).
|
|
func newStoreAway(windows *Windows, log func(string)) *storeAway {
|
|
if log == nil {
|
|
log = func(string) {}
|
|
}
|
|
return &storeAway{windows: windows, log: log, left: storeAwayAtMost, said: map[string]bool{}}
|
|
}
|
|
|
|
// next is how long to wait before fetching again after the attempt-th failure, or zero and why not.
|
|
func (a *storeAway) next(attempt int) (time.Duration, string) {
|
|
a.mu.Lock()
|
|
defer a.mu.Unlock()
|
|
if open := a.windows.Open(time.Now()); len(open) > 0 {
|
|
if a.windowWaited < windowWaitAtMost {
|
|
if key := "window " + open[0].Step; !a.said[key] {
|
|
a.said[key] = true
|
|
a.log(fmt.Sprintf(" waiting the store does not answer while a maintenance window is open "+
|
|
"on this machine (%s): fetched again once it closes, for at most %s (novox/hq issue 291)",
|
|
open[0], windowWaitAtMost))
|
|
}
|
|
a.windowWaited += windowPoll
|
|
return windowPoll, ""
|
|
}
|
|
return 0, fmt.Sprintf("waited %s for the maintenance window open on this machine (%s), and it did not "+
|
|
"close", a.windowWaited.Round(time.Second), open[0])
|
|
}
|
|
if a.left <= 0 {
|
|
return 0, fmt.Sprintf("the store did not answer for the %s this apply waits for it; the next pass "+
|
|
"fetches it again", storeAwayAtMost)
|
|
}
|
|
wait := storeBackoff
|
|
for i := 1; i < attempt && wait < storeBackoffMax; i++ {
|
|
wait *= 2
|
|
}
|
|
wait = min(wait, storeBackoffMax, a.left)
|
|
a.left -= wait
|
|
if !a.said["away"] {
|
|
a.said["away"] = true
|
|
a.log(fmt.Sprintf(" waiting the store does not answer — its maintenance window, perhaps, on "+
|
|
"its own machine: fetched again for at most %s in this apply (novox/hq issue 291)", storeAwayAtMost))
|
|
}
|
|
return wait, ""
|
|
}
|
|
|
|
// fetchMinding is fetch, patient with a store that is away on purpose. A nil away is not patient: a
|
|
// caller outside an apply.
|
|
func fetchMinding(ctx context.Context, source string, away *storeAway) ([]byte, error) {
|
|
for attempt := 1; ; attempt++ {
|
|
body, err := fetch(ctx, source)
|
|
if err == nil || away == nil || !storeAwayErr(err) {
|
|
return body, err
|
|
}
|
|
wait, why := away.next(attempt)
|
|
if wait <= 0 {
|
|
if why != "" && !strings.Contains(err.Error(), why) {
|
|
err = fmt.Errorf("%w — %s", err, why)
|
|
}
|
|
return nil, err
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, err
|
|
case <-time.After(wait):
|
|
}
|
|
}
|
|
}
|