VictoriaMetrics/vendor/google.golang.org/grpc/pickfirst.go

228 lines
6.4 KiB
Go
Raw Normal View History

/*
*
* Copyright 2017 gRPC 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 grpc
import (
2023-07-07 09:05:50 +02:00
"encoding/json"
2019-12-26 18:42:48 +01:00
"errors"
2020-06-25 22:42:41 +02:00
"fmt"
"google.golang.org/grpc/balancer"
"google.golang.org/grpc/connectivity"
2023-07-07 09:05:50 +02:00
"google.golang.org/grpc/internal/envconfig"
"google.golang.org/grpc/internal/grpcrand"
"google.golang.org/grpc/serviceconfig"
)
// PickFirstBalancerName is the name of the pick_first balancer.
const PickFirstBalancerName = "pick_first"
func newPickfirstBuilder() balancer.Builder {
return &pickfirstBuilder{}
}
type pickfirstBuilder struct{}
func (*pickfirstBuilder) Build(cc balancer.ClientConn, opt balancer.BuildOptions) balancer.Balancer {
return &pickfirstBalancer{cc: cc}
}
func (*pickfirstBuilder) Name() string {
return PickFirstBalancerName
}
2023-07-07 09:05:50 +02:00
type pfConfig struct {
serviceconfig.LoadBalancingConfig `json:"-"`
// If set to true, instructs the LB policy to shuffle the order of the list
// of addresses received from the name resolver before attempting to
// connect to them.
ShuffleAddressList bool `json:"shuffleAddressList"`
}
func (*pickfirstBuilder) ParseConfig(js json.RawMessage) (serviceconfig.LoadBalancingConfig, error) {
cfg := &pfConfig{}
if err := json.Unmarshal(js, cfg); err != nil {
return nil, fmt.Errorf("pickfirst: unable to unmarshal LB policy config: %s, error: %v", string(js), err)
}
return cfg, nil
}
type pickfirstBalancer struct {
2022-04-26 14:24:23 +02:00
state connectivity.State
cc balancer.ClientConn
subConn balancer.SubConn
2023-07-07 09:05:50 +02:00
cfg *pfConfig
}
2019-12-26 18:42:48 +01:00
func (b *pickfirstBalancer) ResolverError(err error) {
2020-08-05 10:15:42 +02:00
if logger.V(2) {
2023-02-08 17:55:14 +01:00
logger.Infof("pickfirstBalancer: ResolverError called with error: %v", err)
2019-12-26 18:42:48 +01:00
}
2022-04-26 14:24:23 +02:00
if b.subConn == nil {
b.state = connectivity.TransientFailure
}
if b.state != connectivity.TransientFailure {
// The picker will not change since the balancer does not currently
// report an error.
return
}
b.cc.UpdateState(balancer.State{
ConnectivityState: connectivity.TransientFailure,
Picker: &picker{err: fmt.Errorf("name resolver error: %v", err)},
})
2019-12-26 18:42:48 +01:00
}
2022-04-26 14:24:23 +02:00
func (b *pickfirstBalancer) UpdateClientConnState(state balancer.ClientConnState) error {
2023-07-07 09:05:50 +02:00
addrs := state.ResolverState.Addresses
if len(addrs) == 0 {
2022-04-26 14:24:23 +02:00
// The resolver reported an empty address list. Treat it like an error by
// calling b.ResolverError.
if b.subConn != nil {
// Remove the old subConn. All addresses were removed, so it is no longer
// valid.
b.cc.RemoveSubConn(b.subConn)
b.subConn = nil
}
2019-12-26 18:42:48 +01:00
b.ResolverError(errors.New("produced zero addresses"))
return balancer.ErrBadResolverState
}
2022-04-26 14:24:23 +02:00
2023-07-07 09:05:50 +02:00
if state.BalancerConfig != nil {
cfg, ok := state.BalancerConfig.(*pfConfig)
if !ok {
return fmt.Errorf("pickfirstBalancer: received nil or illegal BalancerConfig (type %T): %v", state.BalancerConfig, state.BalancerConfig)
}
b.cfg = cfg
}
if envconfig.PickFirstLBConfig && b.cfg != nil && b.cfg.ShuffleAddressList {
grpcrand.Shuffle(len(addrs), func(i, j int) { addrs[i], addrs[j] = addrs[j], addrs[i] })
}
2022-04-26 14:24:23 +02:00
if b.subConn != nil {
2023-07-07 09:05:50 +02:00
b.cc.UpdateAddresses(b.subConn, addrs)
2022-04-26 14:24:23 +02:00
return nil
}
2023-07-07 09:05:50 +02:00
subConn, err := b.cc.NewSubConn(addrs, balancer.NewSubConnOptions{})
2022-04-26 14:24:23 +02:00
if err != nil {
if logger.V(2) {
logger.Errorf("pickfirstBalancer: failed to NewSubConn: %v", err)
}
2022-04-26 14:24:23 +02:00
b.state = connectivity.TransientFailure
b.cc.UpdateState(balancer.State{
ConnectivityState: connectivity.TransientFailure,
Picker: &picker{err: fmt.Errorf("error creating connection: %v", err)},
})
return balancer.ErrBadResolverState
}
2022-04-26 14:24:23 +02:00
b.subConn = subConn
b.state = connectivity.Idle
b.cc.UpdateState(balancer.State{
2023-01-11 03:58:34 +01:00
ConnectivityState: connectivity.Connecting,
Picker: &picker{err: balancer.ErrNoSubConnAvailable},
2022-04-26 14:24:23 +02:00
})
b.subConn.Connect()
2019-12-26 18:42:48 +01:00
return nil
}
2022-04-26 14:24:23 +02:00
func (b *pickfirstBalancer) UpdateSubConnState(subConn balancer.SubConn, state balancer.SubConnState) {
2020-08-05 10:15:42 +02:00
if logger.V(2) {
2022-04-26 14:24:23 +02:00
logger.Infof("pickfirstBalancer: UpdateSubConnState: %p, %v", subConn, state)
}
2022-04-26 14:24:23 +02:00
if b.subConn != subConn {
2020-08-05 10:15:42 +02:00
if logger.V(2) {
2022-04-26 14:24:23 +02:00
logger.Infof("pickfirstBalancer: ignored state change because subConn is not recognized")
}
return
}
2022-04-26 14:24:23 +02:00
if state.ConnectivityState == connectivity.Shutdown {
b.subConn = nil
return
}
2022-04-26 14:24:23 +02:00
switch state.ConnectivityState {
2021-09-27 17:01:40 +02:00
case connectivity.Ready:
2022-04-26 14:24:23 +02:00
b.cc.UpdateState(balancer.State{
ConnectivityState: state.ConnectivityState,
Picker: &picker{result: balancer.PickResult{SubConn: subConn}},
})
case connectivity.Connecting:
2023-07-07 09:05:50 +02:00
if b.state == connectivity.TransientFailure {
// We stay in TransientFailure until we are Ready. See A62.
return
}
2022-04-26 14:24:23 +02:00
b.cc.UpdateState(balancer.State{
ConnectivityState: state.ConnectivityState,
Picker: &picker{err: balancer.ErrNoSubConnAvailable},
})
2021-09-27 17:01:40 +02:00
case connectivity.Idle:
2023-07-07 09:05:50 +02:00
if b.state == connectivity.TransientFailure {
// We stay in TransientFailure until we are Ready. Also kick the
// subConn out of Idle into Connecting. See A62.
b.subConn.Connect()
return
}
2022-04-26 14:24:23 +02:00
b.cc.UpdateState(balancer.State{
ConnectivityState: state.ConnectivityState,
Picker: &idlePicker{subConn: subConn},
})
case connectivity.TransientFailure:
2019-12-26 18:42:48 +01:00
b.cc.UpdateState(balancer.State{
2022-04-26 14:24:23 +02:00
ConnectivityState: state.ConnectivityState,
Picker: &picker{err: state.ConnectionError},
2019-12-26 18:42:48 +01:00
})
}
2023-07-07 09:05:50 +02:00
b.state = state.ConnectivityState
}
func (b *pickfirstBalancer) Close() {
}
2021-09-27 17:01:40 +02:00
func (b *pickfirstBalancer) ExitIdle() {
2022-04-26 14:24:23 +02:00
if b.subConn != nil && b.state == connectivity.Idle {
b.subConn.Connect()
2021-09-27 17:01:40 +02:00
}
}
type picker struct {
2019-12-26 18:42:48 +01:00
result balancer.PickResult
err error
}
2022-04-26 14:24:23 +02:00
func (p *picker) Pick(balancer.PickInfo) (balancer.PickResult, error) {
2019-12-26 18:42:48 +01:00
return p.result, p.err
}
2021-09-27 17:01:40 +02:00
// idlePicker is used when the SubConn is IDLE and kicks the SubConn into
// CONNECTING when Pick is called.
type idlePicker struct {
2022-04-26 14:24:23 +02:00
subConn balancer.SubConn
2021-09-27 17:01:40 +02:00
}
2022-04-26 14:24:23 +02:00
func (i *idlePicker) Pick(balancer.PickInfo) (balancer.PickResult, error) {
i.subConn.Connect()
2021-09-27 17:01:40 +02:00
return balancer.PickResult{}, balancer.ErrNoSubConnAvailable
}
func init() {
balancer.Register(newPickfirstBuilder())
}