mirror of
https://github.com/caddyserver/caddy.git
synced 2025-01-27 23:03:37 -05:00
proxy: Added unbuffered request optimization
If only one upstream is defined we don't need to buffer the body. Instead we directly stream the body to the upstream host, which reduces memory usage as well as latency. Furthermore this enables different kinds of HTTP streaming applications like gRPC for instance.
This commit is contained in:
parent
fab3b5bf19
commit
8048e9c3bc
2 changed files with 43 additions and 13 deletions
|
@ -39,6 +39,9 @@ type Upstream interface {
|
||||||
// Gets how long to wait between selecting upstream
|
// Gets how long to wait between selecting upstream
|
||||||
// hosts in the case of cascading failures.
|
// hosts in the case of cascading failures.
|
||||||
GetTryInterval() time.Duration
|
GetTryInterval() time.Duration
|
||||||
|
|
||||||
|
// Gets the number of upstream hosts.
|
||||||
|
GetHostCount() int
|
||||||
}
|
}
|
||||||
|
|
||||||
// UpstreamHostDownFunc can be used to customize how Down behaves.
|
// UpstreamHostDownFunc can be used to customize how Down behaves.
|
||||||
|
@ -94,7 +97,19 @@ func (p Proxy) ServeHTTP(w http.ResponseWriter, r *http.Request) (int, error) {
|
||||||
// outreq is the request that makes a roundtrip to the backend
|
// outreq is the request that makes a roundtrip to the backend
|
||||||
outreq := createUpstreamRequest(r)
|
outreq := createUpstreamRequest(r)
|
||||||
|
|
||||||
// record and replace outreq body
|
// If we have more than one upstream host defined and if retrying is enabled
|
||||||
|
// by setting try_duration to a non-zero value, caddy will try to
|
||||||
|
// retry the request at a different host if the first one failed.
|
||||||
|
//
|
||||||
|
// This requires us to possibly rewind and replay the request body though,
|
||||||
|
// which in turn requires us to buffer the request body first.
|
||||||
|
//
|
||||||
|
// An unbuffered request is usually preferrable, because it reduces latency
|
||||||
|
// as well as memory usage. Furthermore it enables different kinds of
|
||||||
|
// HTTP streaming applications like gRPC for instance.
|
||||||
|
requiresBuffering := upstream.GetHostCount() > 1 && upstream.GetTryDuration() != 0
|
||||||
|
|
||||||
|
if requiresBuffering {
|
||||||
body, err := newBufferedBody(outreq.Body)
|
body, err := newBufferedBody(outreq.Body)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return http.StatusBadRequest, errors.New("failed to read downstream request body")
|
return http.StatusBadRequest, errors.New("failed to read downstream request body")
|
||||||
|
@ -102,6 +117,7 @@ func (p Proxy) ServeHTTP(w http.ResponseWriter, r *http.Request) (int, error) {
|
||||||
if body != nil {
|
if body != nil {
|
||||||
outreq.Body = body
|
outreq.Body = body
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// The keepRetrying function will return true if we should
|
// The keepRetrying function will return true if we should
|
||||||
// loop and try to select another host, or false if we
|
// loop and try to select another host, or false if we
|
||||||
|
@ -173,15 +189,25 @@ func (p Proxy) ServeHTTP(w http.ResponseWriter, r *http.Request) (int, error) {
|
||||||
downHeaderUpdateFn = createRespHeaderUpdateFn(host.DownstreamHeaders, replacer)
|
downHeaderUpdateFn = createRespHeaderUpdateFn(host.DownstreamHeaders, replacer)
|
||||||
}
|
}
|
||||||
|
|
||||||
// rewind request body to its beginning
|
// Before we retry the request we have to make sure
|
||||||
if err := body.rewind(); err != nil {
|
// that the body is rewound to it's beginning.
|
||||||
|
if bb, ok := outreq.Body.(*bufferedBody); ok {
|
||||||
|
if err := bb.rewind(); err != nil {
|
||||||
return http.StatusInternalServerError, errors.New("unable to rewind downstream request body")
|
return http.StatusInternalServerError, errors.New("unable to rewind downstream request body")
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// tell the proxy to serve the request
|
// tell the proxy to serve the request
|
||||||
|
//
|
||||||
|
// NOTE:
|
||||||
|
// The call to proxy.ServeHTTP can theoretically panic.
|
||||||
|
// To prevent host.Conns from getting out-of-sync we thus have to
|
||||||
|
// make sure that it's _always_ correctly decremented afterwards.
|
||||||
|
func() {
|
||||||
atomic.AddInt64(&host.Conns, 1)
|
atomic.AddInt64(&host.Conns, 1)
|
||||||
|
defer atomic.AddInt64(&host.Conns, -1)
|
||||||
backendErr = proxy.ServeHTTP(w, outreq, downHeaderUpdateFn)
|
backendErr = proxy.ServeHTTP(w, outreq, downHeaderUpdateFn)
|
||||||
atomic.AddInt64(&host.Conns, -1)
|
}()
|
||||||
|
|
||||||
// if no errors, we're done here
|
// if no errors, we're done here
|
||||||
if backendErr == nil {
|
if backendErr == nil {
|
||||||
|
|
|
@ -423,6 +423,10 @@ func (u *staticUpstream) GetTryInterval() time.Duration {
|
||||||
return u.TryInterval
|
return u.TryInterval
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (u *staticUpstream) GetHostCount() int {
|
||||||
|
return len(u.Hosts)
|
||||||
|
}
|
||||||
|
|
||||||
// RegisterPolicy adds a custom policy to the proxy.
|
// RegisterPolicy adds a custom policy to the proxy.
|
||||||
func RegisterPolicy(name string, policy func() Policy) {
|
func RegisterPolicy(name string, policy func() Policy) {
|
||||||
supportedPolicies[name] = policy
|
supportedPolicies[name] = policy
|
||||||
|
|
Loading…
Add table
Reference in a new issue