filer.sync: back off on transient upload errors (#9777)

A destination volume server that hits its idle deadline while reading a large upload body under load returns 400 "read tcp ...: i/o timeout". fetchAndWrite retried that on the flat ~1s RetryUntil backoff, hammering the already-overloaded destination. Route i/o timeout, connection reset, broken pipe and net.Error timeouts through the same escalating 10s-2min backoff already used for EOF so it can recover.
This commit is contained in:
Chris Lu
2026-06-01 12:18:17 -07:00
committed by GitHub
parent f9ee49b03e
commit 2e3fabbf24
2 changed files with 86 additions and 17 deletions
+42 -14
View File
@@ -6,6 +6,7 @@ import (
"errors"
"fmt"
"io"
"net"
"net/http"
"os"
"path/filepath"
@@ -251,7 +252,7 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
fs.activeTransfers.Store(sourceChunk.GetFileIdString(), transferStatus)
defer fs.activeTransfers.Delete(sourceChunk.GetFileIdString())
eofBackoff := time.Duration(0)
transientBackoff := time.Duration(0)
var partialData []byte
var savedFilename string
var savedHeader http.Header
@@ -335,7 +336,7 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
return fmt.Errorf("upload result: %v", uploadResult.Error)
}
eofBackoff = 0
transientBackoff = 0
fileId = currentFileId
return nil
}, func(retryErr error) (shouldContinue bool) {
@@ -354,15 +355,15 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
transferStatus.mu.Lock()
transferStatus.LastErr = retryErr.Error()
transferStatus.mu.Unlock()
if isEofError(retryErr) {
eofBackoff = nextEofBackoff(eofBackoff)
if isRetryableNetworkError(retryErr) {
transientBackoff = nextTransientBackoff(transientBackoff)
transferStatus.mu.Lock()
transferStatus.BytesReceived = int64(len(partialData))
transferStatus.Status = fmt.Sprintf("waiting %v", eofBackoff)
transferStatus.Status = fmt.Sprintf("waiting %v", transientBackoff)
transferStatus.mu.Unlock()
glog.V(0).Infof("source connection interrupted while replicating %s for %s (%d bytes received so far), backing off %v: %v",
sourceChunk.GetFileIdString(), path, len(partialData), eofBackoff, retryErr)
time.Sleep(eofBackoff)
glog.V(0).Infof("connection interrupted while replicating %s for %s (%d bytes received so far), backing off %v: %v",
sourceChunk.GetFileIdString(), path, len(partialData), transientBackoff, retryErr)
time.Sleep(transientBackoff)
transferStatus.mu.Lock()
transferStatus.Status = "downloading"
transferStatus.mu.Unlock()
@@ -378,17 +379,19 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
return fileId, nil
}
const maxEofBackoff = 2 * time.Minute
const maxTransientBackoff = 2 * time.Minute
// nextEofBackoff returns the next backoff duration for unexpected EOF errors.
// It starts at 10s, doubles each time, and caps at 2 minutes.
func nextEofBackoff(current time.Duration) time.Duration {
// nextTransientBackoff returns the next backoff duration for a transient
// network failure. It starts at 10s, doubles each time, and caps at 2 minutes
// so an overloaded destination can recover instead of being hammered every
// second by the surrounding RetryUntil loop.
func nextTransientBackoff(current time.Duration) time.Duration {
if current < 10*time.Second {
return 10 * time.Second
}
current *= 2
if current > maxEofBackoff {
current = maxEofBackoff
if current > maxTransientBackoff {
current = maxTransientBackoff
}
return current
}
@@ -400,6 +403,31 @@ func isEofError(err error) bool {
return errors.Is(err, io.ErrUnexpectedEOF) || errors.Is(err, io.EOF)
}
// isRetryableNetworkError reports whether err is a transient network failure
// worth a backoff-and-retry: an interrupted read (EOF), a timeout (e.g. the
// destination volume server hitting its idle deadline while reading a large
// upload body under load), or a reset/broken connection. The volume server
// returns the timeout as a JSON error string, so match on text in addition to
// the net.Error interface.
func isRetryableNetworkError(err error) bool {
if err == nil {
return false
}
if isEofError(err) {
return true
}
var netErr net.Error
if errors.As(err, &netErr) && netErr.Timeout() {
return true
}
// lower-case so we also catch capitalized variants from other OSes,
// libraries, or custom error wrappers
msg := strings.ToLower(err.Error())
return strings.Contains(msg, "i/o timeout") ||
strings.Contains(msg, "connection reset") ||
strings.Contains(msg, "broken pipe")
}
// errChunkSizeMismatch is a permanent (non-retriable) replication failure.
var errChunkSizeMismatch = errors.New("chunk size mismatch")
@@ -2,6 +2,8 @@ package filersink
import (
"errors"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
@@ -164,9 +166,9 @@ func TestValidateReplicatedChunkSize(t *testing.T) {
}
// End-to-end regression :
// a source volume that responds 200 OK with Content-Length: 0
// for a chunk that filer metadata claims is 5171 bytes must be rejected
// by fetchAndWrite with a (non-retriable) size mismatch error,
// a source volume that responds 200 OK with Content-Length: 0
// for a chunk that filer metadata claims is 5171 bytes must be rejected
// by fetchAndWrite with a (non-retriable) size mismatch error,
// instead of being silently propagated to the destination as a 0-byte needle.
func TestFetchAndWriteRejectsZeroByteSource(t *testing.T) {
const fid = "74,047d16a94aa581"
@@ -242,6 +244,45 @@ func TestFetchAndWriteRejectsZeroByteSource(t *testing.T) {
}
}
type timeoutErr struct{}
func (timeoutErr) Error() string { return "synthetic timeout" }
func (timeoutErr) Timeout() bool { return true }
func (timeoutErr) Temporary() bool { return true }
// A transient network failure (interrupted read, idle-deadline timeout while
// the destination reads the upload body, reset/broken pipe) must route through
// the escalating backoff so an overloaded destination can recover instead of
// being hammered. The volume server returns its idle timeout as a JSON error
// string, so the text path matters as much as the net.Error interface.
func TestIsRetryableNetworkError(t *testing.T) {
tests := []struct {
name string
err error
want bool
}{
{"nil", nil, false},
{"eof", io.EOF, true},
{"unexpected eof", io.ErrUnexpectedEOF, true},
{"volume idle timeout json", fmt.Errorf("upload result: read tcp 10.0.0.1:8082->10.0.0.1:54848: i/o timeout"), true},
{"volume idle timeout capitalized", fmt.Errorf("Upload result: read tcp 10.0.0.1:8082->10.0.0.1:54848: I/O timeout"), true},
{"connection reset", fmt.Errorf("upload data: write tcp ...: connection reset by peer"), true},
{"connection reset capitalized", fmt.Errorf("Connection reset by peer"), true},
{"broken pipe", fmt.Errorf("broken pipe"), true},
{"broken pipe capitalized", fmt.Errorf("Broken pipe"), true},
{"net.Error timeout", fmt.Errorf("dial: %w", timeoutErr{}), true},
{"size mismatch is permanent", errChunkSizeMismatch, false},
{"unrelated error", errors.New("not found"), false},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
if got := isRetryableNetworkError(tc.err); got != tc.want {
t.Fatalf("isRetryableNetworkError(%v) = %v, want %v", tc.err, got, tc.want)
}
})
}
}
// Lock in that the errChunkSizeMismatch sentinel survives the wrap in
// replicateOneChunk + pass-through in util.Retry, so filer_sink.go's
// errors.Is check actually fires.