diff --git a/weed/util/http/client/http_client.go b/weed/util/http/client/http_client.go index 172566c3b..305f78c31 100644 --- a/weed/util/http/client/http_client.go +++ b/weed/util/http/client/http_client.go @@ -10,6 +10,7 @@ import ( "os" "strings" "sync" + "time" "google.golang.org/grpc/credentials/tls/certprovider" @@ -22,6 +23,17 @@ var ( loadSecurityConfigOnce sync.Once ) +// Intra-cluster peers (volume/filer/master) answer headers within milliseconds. +// A peer that is TCP-reachable but not answering -- a volume server still loading +// after a restart, or a stale keep-alive to a container that returned on a new IP +// -- otherwise blocks a chunk read or a replicated write forever. The response +// timeout bounds that wait so the read fails over to another replica and the +// write fails fast to retry; the idle timeout evicts sockets to a departed server. +const ( + responseHeaderTimeout = 30 * time.Second + idleConnTimeout = 90 * time.Second +) + type HTTPClient struct { Client *http.Client Transport *http.Transport @@ -166,7 +178,9 @@ func NewHttpClient(clientName ClientName, opts ...HttpClientOpt) (*HTTPClient, e MaxIdleConnsPerHost: 1024, TLSClientConfig: tlsConfig, // Bind outbound HTTP to the -ip.bind source address. - DialContext: util.OutboundDialContext, + DialContext: util.OutboundDialContext, + ResponseHeaderTimeout: responseHeaderTimeout, + IdleConnTimeout: idleConnTimeout, } httpClient.Client = &http.Client{ Transport: httpClient.Transport, @@ -282,7 +296,9 @@ func NewHttpClientWithTLS(certFile, keyFile, caFile string, insecureSkipVerify b MaxIdleConnsPerHost: 1024, TLSClientConfig: tlsConfig, // Bind outbound HTTP to the -ip.bind source address. - DialContext: util.OutboundDialContext, + DialContext: util.OutboundDialContext, + ResponseHeaderTimeout: responseHeaderTimeout, + IdleConnTimeout: idleConnTimeout, } httpClient.Client = &http.Client{ Transport: httpClient.Transport, diff --git a/weed/util/http/client/http_client_opt.go b/weed/util/http/client/http_client_opt.go index 1ff9d533d..8b5c1e55c 100644 --- a/weed/util/http/client/http_client_opt.go +++ b/weed/util/http/client/http_client_opt.go @@ -16,3 +16,14 @@ func AddDialContext(httpClient *HTTPClient) { httpClient.Transport.DialContext = dialContext httpClient.Client.Transport = httpClient.Transport } + +// WithResponseHeaderTimeout overrides the transport's default response-header +// timeout. Mainly for tests that need a short deadline against an unresponsive +// peer without mutating the package default. +func WithResponseHeaderTimeout(timeout time.Duration) HttpClientOpt { + return func(httpClient *HTTPClient) { + if httpClient.Transport != nil { + httpClient.Transport.ResponseHeaderTimeout = timeout + } + } +} diff --git a/weed/util/http/client/http_client_test.go b/weed/util/http/client/http_client_test.go new file mode 100644 index 000000000..546e1c518 --- /dev/null +++ b/weed/util/http/client/http_client_test.go @@ -0,0 +1,127 @@ +package client + +import ( + "context" + "errors" + "net" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestNewHttpClientSetsResponseTimeouts(t *testing.T) { + c, err := NewHttpClient(Client) + if err != nil { + t.Fatalf("NewHttpClient: %v", err) + } + if got := c.Transport.ResponseHeaderTimeout; got != responseHeaderTimeout { + t.Errorf("ResponseHeaderTimeout = %v, want %v", got, responseHeaderTimeout) + } + if got := c.Transport.IdleConnTimeout; got != idleConnTimeout { + t.Errorf("IdleConnTimeout = %v, want %v", got, idleConnTimeout) + } + // AddDialContext must keep the response timeouts intact. + c2, err := NewHttpClient(Client, AddDialContext) + if err != nil { + t.Fatalf("NewHttpClient(AddDialContext): %v", err) + } + if got := c2.Transport.ResponseHeaderTimeout; got != responseHeaderTimeout { + t.Errorf("AddDialContext dropped ResponseHeaderTimeout: got %v", got) + } + // WithResponseHeaderTimeout overrides the default. + c3, err := NewHttpClient(Client, WithResponseHeaderTimeout(time.Second)) + if err != nil { + t.Fatalf("NewHttpClient(WithResponseHeaderTimeout): %v", err) + } + if got := c3.Transport.ResponseHeaderTimeout; got != time.Second { + t.Errorf("WithResponseHeaderTimeout not applied: got %v", got) + } + // The TLS constructor must carry the same defaults. + cTLS, err := NewHttpClientWithTLS("", "", "", true) + if err != nil { + t.Fatalf("NewHttpClientWithTLS: %v", err) + } + if got := cTLS.Transport.ResponseHeaderTimeout; got != responseHeaderTimeout { + t.Errorf("NewHttpClientWithTLS ResponseHeaderTimeout = %v, want %v", got, responseHeaderTimeout) + } + if got := cTLS.Transport.IdleConnTimeout; got != idleConnTimeout { + t.Errorf("NewHttpClientWithTLS IdleConnTimeout = %v, want %v", got, idleConnTimeout) + } +} + +// A peer that is TCP-reachable but never answers -- a volume server still +// loading its volumes after a restart, or a stale keep-alive to a container +// that came back on a new IP -- must not block a read or a replicated write +// forever. ResponseHeaderTimeout turns that into a prompt timeout error so the +// caller can fail over to another replica or retry. +func TestResponseHeaderTimeoutUnblocksUnresponsivePeer(t *testing.T) { + server, release := newUnresponsiveServer(t, false) + defer server.Close() + defer close(release) + + c, err := NewHttpClient(Client, WithResponseHeaderTimeout(200*time.Millisecond)) + if err != nil { + t.Fatalf("NewHttpClient: %v", err) + } + assertTimesOut(t, c, server.URL) +} + +// Same guarantee over the TLS client path (NewHttpClientWithTLS). +func TestResponseHeaderTimeoutUnblocksUnresponsivePeerTLS(t *testing.T) { + server, release := newUnresponsiveServer(t, true) + defer server.Close() + defer close(release) + + c, err := NewHttpClientWithTLS("", "", "", true, WithResponseHeaderTimeout(200*time.Millisecond)) + if err != nil { + t.Fatalf("NewHttpClientWithTLS: %v", err) + } + assertTimesOut(t, c, server.URL) +} + +// newUnresponsiveServer returns a server whose handler blocks until the returned +// channel is closed, so it accepts the connection but never sends response +// headers. Callers must close the channel before server.Close() -- defers run +// LIFO, so register the Close defer first. +func newUnresponsiveServer(t *testing.T, tls bool) (*httptest.Server, chan struct{}) { + t.Helper() + release := make(chan struct{}) + handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + <-release + }) + if tls { + return httptest.NewTLSServer(handler), release + } + return httptest.NewServer(handler), release +} + +func assertTimesOut(t *testing.T, c *HTTPClient, url string) { + t.Helper() + req, err := http.NewRequestWithContext(context.Background(), http.MethodGet, url, nil) + if err != nil { + t.Fatalf("new request: %v", err) + } + + done := make(chan error, 1) + go func() { + resp, doErr := c.Do(req) + if resp != nil { + resp.Body.Close() + } + done <- doErr + }() + + select { + case doErr := <-done: + if doErr == nil { + t.Fatal("expected a timeout error from an unresponsive peer, got nil") + } + var netErr net.Error + if !errors.As(doErr, &netErr) || !netErr.Timeout() { + t.Fatalf("expected a timeout error, got %v", doErr) + } + case <-time.After(5 * time.Second): + t.Fatal("Do() blocked well past the response header timeout") + } +}