2019-09-04 01:56:09 +03:00
|
|
|
// Copyright 2015 Matthew Holt and The Caddy Authors
|
|
|
|
//
|
|
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
|
|
// you may not use this file except in compliance with the License.
|
|
|
|
// You may obtain a copy of the License at
|
|
|
|
//
|
|
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
//
|
|
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
// See the License for the specific language governing permissions and
|
|
|
|
// limitations under the License.
|
|
|
|
|
|
|
|
package reverseproxy
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
2019-09-14 22:25:26 +03:00
|
|
|
"net"
|
|
|
|
"strings"
|
2019-09-04 01:56:09 +03:00
|
|
|
"sync/atomic"
|
|
|
|
|
|
|
|
"github.com/caddyserver/caddy/v2"
|
|
|
|
)
|
|
|
|
|
|
|
|
// Host represents a remote host which can be proxied to.
|
|
|
|
// Its methods must be safe for concurrent use.
|
|
|
|
type Host interface {
|
|
|
|
// NumRequests returns the numnber of requests
|
|
|
|
// currently in process with the host.
|
|
|
|
NumRequests() int
|
|
|
|
|
|
|
|
// Fails returns the count of recent failures.
|
|
|
|
Fails() int
|
|
|
|
|
|
|
|
// Unhealthy returns true if the backend is unhealthy.
|
|
|
|
Unhealthy() bool
|
|
|
|
|
2019-09-10 06:44:58 +03:00
|
|
|
// CountRequest atomically counts the given number of
|
|
|
|
// requests as currently in process with the host. The
|
|
|
|
// count should not go below 0.
|
2019-09-04 01:56:09 +03:00
|
|
|
CountRequest(int) error
|
|
|
|
|
2019-09-10 06:44:58 +03:00
|
|
|
// CountFail atomically counts the given number of
|
|
|
|
// failures with the host. The count should not go
|
|
|
|
// below 0.
|
2019-09-04 01:56:09 +03:00
|
|
|
CountFail(int) error
|
|
|
|
|
2019-09-10 06:44:58 +03:00
|
|
|
// SetHealthy atomically marks the host as either
|
|
|
|
// healthy (true) or unhealthy (false). If the given
|
|
|
|
// status is the same, this should be a no-op and
|
|
|
|
// return false. It returns true if the status was
|
|
|
|
// changed; i.e. if it is now different from before.
|
2019-09-04 01:56:09 +03:00
|
|
|
SetHealthy(bool) (bool, error)
|
|
|
|
}
|
|
|
|
|
|
|
|
// UpstreamPool is a collection of upstreams.
|
|
|
|
type UpstreamPool []*Upstream
|
|
|
|
|
|
|
|
// Upstream bridges this proxy's configuration to the
|
|
|
|
// state of the backend host it is correlated with.
|
|
|
|
type Upstream struct {
|
|
|
|
Host `json:"-"`
|
|
|
|
|
2019-09-05 22:14:39 +03:00
|
|
|
Dial string `json:"dial,omitempty"`
|
2019-09-04 01:56:09 +03:00
|
|
|
MaxRequests int `json:"max_requests,omitempty"`
|
|
|
|
|
|
|
|
// TODO: This could be really useful, to bind requests
|
|
|
|
// with certain properties to specific backends
|
|
|
|
// HeaderAffinity string
|
|
|
|
// IPAffinity string
|
|
|
|
|
|
|
|
healthCheckPolicy *PassiveHealthChecks
|
2019-09-04 04:06:54 +03:00
|
|
|
cb CircuitBreaker
|
2019-09-04 01:56:09 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
// Available returns true if the remote host
|
2019-09-04 04:06:54 +03:00
|
|
|
// is available to receive requests. This is
|
|
|
|
// the method that should be used by selection
|
|
|
|
// policies, etc. to determine if a backend
|
|
|
|
// should be able to be sent a request.
|
2019-09-04 01:56:09 +03:00
|
|
|
func (u *Upstream) Available() bool {
|
|
|
|
return u.Healthy() && !u.Full()
|
|
|
|
}
|
|
|
|
|
|
|
|
// Healthy returns true if the remote host
|
|
|
|
// is currently known to be healthy or "up".
|
2019-09-04 04:06:54 +03:00
|
|
|
// It consults the circuit breaker, if any.
|
2019-09-04 01:56:09 +03:00
|
|
|
func (u *Upstream) Healthy() bool {
|
|
|
|
healthy := !u.Host.Unhealthy()
|
|
|
|
if healthy && u.healthCheckPolicy != nil {
|
|
|
|
healthy = u.Host.Fails() < u.healthCheckPolicy.MaxFails
|
|
|
|
}
|
2019-09-04 04:06:54 +03:00
|
|
|
if healthy && u.cb != nil {
|
|
|
|
healthy = u.cb.OK()
|
|
|
|
}
|
2019-09-04 01:56:09 +03:00
|
|
|
return healthy
|
|
|
|
}
|
|
|
|
|
|
|
|
// Full returns true if the remote host
|
|
|
|
// cannot receive more requests at this time.
|
|
|
|
func (u *Upstream) Full() bool {
|
|
|
|
return u.MaxRequests > 0 && u.Host.NumRequests() >= u.MaxRequests
|
|
|
|
}
|
|
|
|
|
|
|
|
// upstreamHost is the basic, in-memory representation
|
|
|
|
// of the state of a remote host. It implements the
|
|
|
|
// Host interface.
|
|
|
|
type upstreamHost struct {
|
|
|
|
numRequests int64 // must be first field to be 64-bit aligned on 32-bit systems (see https://golang.org/pkg/sync/atomic/#pkg-note-BUG)
|
|
|
|
fails int64
|
|
|
|
unhealthy int32
|
|
|
|
}
|
|
|
|
|
|
|
|
// NumRequests returns the number of active requests to the upstream.
|
|
|
|
func (uh *upstreamHost) NumRequests() int {
|
|
|
|
return int(atomic.LoadInt64(&uh.numRequests))
|
|
|
|
}
|
|
|
|
|
|
|
|
// Fails returns the number of recent failures with the upstream.
|
|
|
|
func (uh *upstreamHost) Fails() int {
|
|
|
|
return int(atomic.LoadInt64(&uh.fails))
|
|
|
|
}
|
|
|
|
|
|
|
|
// Unhealthy returns whether the upstream is healthy.
|
|
|
|
func (uh *upstreamHost) Unhealthy() bool {
|
|
|
|
return atomic.LoadInt32(&uh.unhealthy) == 1
|
|
|
|
}
|
|
|
|
|
|
|
|
// CountRequest mutates the active request count by
|
|
|
|
// delta. It returns an error if the adjustment fails.
|
|
|
|
func (uh *upstreamHost) CountRequest(delta int) error {
|
|
|
|
result := atomic.AddInt64(&uh.numRequests, int64(delta))
|
|
|
|
if result < 0 {
|
|
|
|
return fmt.Errorf("count below 0: %d", result)
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// CountFail mutates the recent failures count by
|
|
|
|
// delta. It returns an error if the adjustment fails.
|
|
|
|
func (uh *upstreamHost) CountFail(delta int) error {
|
|
|
|
result := atomic.AddInt64(&uh.fails, int64(delta))
|
|
|
|
if result < 0 {
|
|
|
|
return fmt.Errorf("count below 0: %d", result)
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// SetHealthy sets the upstream has healthy or unhealthy
|
2019-10-11 23:25:39 +03:00
|
|
|
// and returns true if the new value is different.
|
2019-09-04 01:56:09 +03:00
|
|
|
func (uh *upstreamHost) SetHealthy(healthy bool) (bool, error) {
|
|
|
|
var unhealthy, compare int32 = 1, 0
|
|
|
|
if healthy {
|
|
|
|
unhealthy, compare = 0, 1
|
|
|
|
}
|
|
|
|
swapped := atomic.CompareAndSwapInt32(&uh.unhealthy, compare, unhealthy)
|
|
|
|
return swapped, nil
|
|
|
|
}
|
|
|
|
|
2019-09-05 22:14:39 +03:00
|
|
|
// DialInfo contains information needed to dial a
|
|
|
|
// connection to an upstream host. This information
|
|
|
|
// may be different than that which is represented
|
|
|
|
// in a URL (for example, unix sockets don't have
|
|
|
|
// a host that can be represented in a URL, but
|
|
|
|
// they certainly have a network name and address).
|
|
|
|
type DialInfo struct {
|
2019-10-11 23:25:39 +03:00
|
|
|
// Upstream is the Upstream associated with
|
|
|
|
// this DialInfo. It may be nil.
|
|
|
|
Upstream *Upstream
|
|
|
|
|
|
|
|
// The network to use. This should be one of
|
|
|
|
// the values that is accepted by net.Dial:
|
2019-09-05 22:14:39 +03:00
|
|
|
// https://golang.org/pkg/net/#Dial
|
|
|
|
Network string
|
|
|
|
|
|
|
|
// The address to dial. Follows the same
|
|
|
|
// semantics and rules as net.Dial.
|
|
|
|
Address string
|
2019-09-14 22:25:26 +03:00
|
|
|
|
2019-10-11 23:25:39 +03:00
|
|
|
// Host and Port are components of Address.
|
2019-09-14 22:25:26 +03:00
|
|
|
Host, Port string
|
|
|
|
}
|
|
|
|
|
2019-09-05 22:14:39 +03:00
|
|
|
// String returns the Caddy network address form
|
|
|
|
// by joining the network and address with a
|
|
|
|
// forward slash.
|
|
|
|
func (di DialInfo) String() string {
|
2019-10-11 23:25:39 +03:00
|
|
|
return caddy.JoinNetworkAddress(di.Network, di.Host, di.Port)
|
|
|
|
}
|
|
|
|
|
|
|
|
// fillDialInfo returns a filled DialInfo for the given upstream, using
|
|
|
|
// the given Replacer. Note that the returned value is not a pointer.
|
|
|
|
func fillDialInfo(upstream *Upstream, repl caddy.Replacer) (DialInfo, error) {
|
|
|
|
dial := repl.ReplaceAll(upstream.Dial, "")
|
|
|
|
netw, addrs, err := caddy.ParseNetworkAddress(dial)
|
|
|
|
if err != nil {
|
|
|
|
return DialInfo{}, fmt.Errorf("upstream %s: invalid dial address %s: %v", upstream.Dial, dial, err)
|
|
|
|
}
|
|
|
|
if len(addrs) != 1 {
|
|
|
|
return DialInfo{}, fmt.Errorf("upstream %s: dial address must represent precisely one socket: %s represents %d",
|
|
|
|
upstream.Dial, dial, len(addrs))
|
|
|
|
}
|
|
|
|
var dialHost, dialPort string
|
|
|
|
if !strings.Contains(netw, "unix") {
|
|
|
|
dialHost, dialPort, err = net.SplitHostPort(addrs[0])
|
|
|
|
if err != nil {
|
|
|
|
dialHost = addrs[0] // assume there was no port
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return DialInfo{
|
|
|
|
Upstream: upstream,
|
|
|
|
Network: netw,
|
|
|
|
Address: addrs[0],
|
|
|
|
Host: dialHost,
|
|
|
|
Port: dialPort,
|
|
|
|
}, nil
|
2019-09-05 22:14:39 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
// DialInfoCtxKey is used to store a DialInfo
|
|
|
|
// in a context.Context.
|
|
|
|
const DialInfoCtxKey = caddy.CtxKey("dial_info")
|
|
|
|
|
2019-09-04 01:56:09 +03:00
|
|
|
// hosts is the global repository for hosts that are
|
|
|
|
// currently in use by active configuration(s). This
|
|
|
|
// allows the state of remote hosts to be preserved
|
|
|
|
// through config reloads.
|
|
|
|
var hosts = caddy.NewUsagePool()
|