grpc-go/clientconn.go

732 строки
20 KiB
Go
Исходник Обычный вид История

2015-02-06 04:14:05 +03:00
/*
*
* Copyright 2014, Google Inc.
* All rights reserved.
*
* Redistribution and use in source and binary forms, with or without
* modification, are permitted provided that the following conditions are
* met:
*
* * Redistributions of source code must retain the above copyright
* notice, this list of conditions and the following disclaimer.
* * Redistributions in binary form must reproduce the above
* copyright notice, this list of conditions and the following disclaimer
* in the documentation and/or other materials provided with the
* distribution.
* * Neither the name of Google Inc. nor the names of its
* contributors may be used to endorse or promote products derived from
* this software without specific prior written permission.
*
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
* OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
* SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT
* LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,
* DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY
* THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
*
*/
package grpc
2015-02-06 04:14:05 +03:00
import (
"errors"
2015-08-01 05:00:43 +03:00
"fmt"
2015-04-17 23:50:18 +03:00
"net"
"strings"
2015-02-06 04:14:05 +03:00
"sync"
"time"
"golang.org/x/net/context"
"golang.org/x/net/trace"
2016-05-07 01:47:09 +03:00
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials"
2015-05-09 12:43:59 +03:00
"google.golang.org/grpc/grpclog"
"google.golang.org/grpc/transport"
2015-02-06 04:14:05 +03:00
)
var (
2016-05-25 21:28:45 +03:00
// ErrClientConnClosing indicates that the operation is illegal because
// the ClientConn is closing.
ErrClientConnClosing = errors.New("grpc: the client connection is closing")
2016-06-06 22:08:11 +03:00
// ErrClientConnTimeout indicates that the ClientConn cannot establish the
// underlying connections within the specified timeout.
ErrClientConnTimeout = errors.New("grpc: timed out when dialing")
2016-05-25 21:28:45 +03:00
// errNoTransportSecurity indicates that there is no transport security
2016-03-17 02:40:16 +03:00
// being set for ClientConn. Users should either set one or explicitly
2015-08-28 03:21:52 +03:00
// call WithInsecure DialOption to disable security.
2016-05-25 21:28:45 +03:00
errNoTransportSecurity = errors.New("grpc: no transport security set (use grpc.WithInsecure() explicitly or set credentials)")
// errTransportCredentialsMissing indicates that users want to transmit security
// information (e.g., oauth2 token) which requires secure connection on an insecure
// connection.
errTransportCredentialsMissing = errors.New("grpc: the credentials require transport level security (use grpc.WithTransportCredentials() to set)")
// errCredentialsConflict indicates that grpc.WithTransportCredentials()
// and grpc.WithInsecure() are both called for a connection.
errCredentialsConflict = errors.New("grpc: transport credentials are set for an insecure connection (grpc.WithTransportCredentials() and grpc.WithInsecure() are both called)")
2016-05-25 21:28:45 +03:00
// errNetworkIP indicates that the connection is down due to some network I/O error.
errNetworkIO = errors.New("grpc: failed with network I/O error")
// errConnDrain indicates that the connection starts to be drained and does not accept any new RPCs.
errConnDrain = errors.New("grpc: the connection is drained")
// errConnClosing indicates that the connection is closing.
errConnClosing = errors.New("grpc: the connection is closing")
2016-06-06 22:16:33 +03:00
errNoAddr = errors.New("grpc: there is no address available to dial")
2015-07-28 21:12:07 +03:00
// minimum time to give a connection to complete
minConnectTimeout = 20 * time.Second
)
// dialOptions configure a Dial call. dialOptions are set by the DialOption
// values passed to Dial.
type dialOptions struct {
2015-08-28 03:21:52 +03:00
codec Codec
cp Compressor
dc Decompressor
bs backoffStrategy
2016-05-07 01:47:09 +03:00
balancer Balancer
2015-08-28 03:21:52 +03:00
block bool
insecure bool
2016-06-06 22:08:11 +03:00
timeout time.Duration
2015-08-28 03:21:52 +03:00
copts transport.ConnectOptions
}
2015-03-04 04:08:39 +03:00
// DialOption configures how we set up the connection.
type DialOption func(*dialOptions)
// WithCodec returns a DialOption which sets a codec for message marshaling and unmarshaling.
2015-04-02 00:22:53 +03:00
func WithCodec(c Codec) DialOption {
return func(o *dialOptions) {
2015-04-02 00:22:53 +03:00
o.codec = c
}
}
2015-02-06 04:14:05 +03:00
2016-01-25 22:18:41 +03:00
// WithCompressor returns a DialOption which sets a CompressorGenerator for generating message
// compressor.
func WithCompressor(cp Compressor) DialOption {
2016-01-23 05:21:41 +03:00
return func(o *dialOptions) {
o.cp = cp
2016-01-23 05:21:41 +03:00
}
}
2016-01-25 22:18:41 +03:00
// WithDecompressor returns a DialOption which sets a DecompressorGenerator for generating
// message decompressor.
func WithDecompressor(dc Decompressor) DialOption {
2016-01-23 05:21:41 +03:00
return func(o *dialOptions) {
o.dc = dc
2016-01-23 05:21:41 +03:00
}
}
2016-05-13 05:19:14 +03:00
// WithBalancer returns a DialOption which sets a load balancer.
func WithBalancer(b Balancer) DialOption {
return func(o *dialOptions) {
o.balancer = b
}
}
// WithBackoffMaxDelay configures the dialer to use the provided maximum delay
// when backing off after failed connection attempts.
func WithBackoffMaxDelay(md time.Duration) DialOption {
return WithBackoffConfig(BackoffConfig{MaxDelay: md})
}
// WithBackoffConfig configures the dialer to use the provided backoff
// parameters after connection failures.
//
// Use WithBackoffMaxDelay until more parameters on BackoffConfig are opened up
// for use.
func WithBackoffConfig(b BackoffConfig) DialOption {
// Set defaults to ensure that provided BackoffConfig is valid and
// unexported fields get default values.
setDefaults(&b)
return withBackoff(b)
}
// withBackoff sets the backoff strategy used for retries after a
// failed connection attempt.
//
// This can be exported if arbitrary backoff strategies are allowed by gRPC.
func withBackoff(bs backoffStrategy) DialOption {
return func(o *dialOptions) {
o.bs = bs
}
}
// WithBlock returns a DialOption which makes caller of Dial blocks until the underlying
2015-06-05 01:47:02 +03:00
// connection is up. Without this, Dial returns immediately and connecting the server
// happens in background.
func WithBlock() DialOption {
return func(o *dialOptions) {
o.block = true
}
}
2016-01-08 01:18:20 +03:00
// WithInsecure returns a DialOption which disables transport security for this ClientConn.
// Note that transport security is required unless WithInsecure is set.
2015-08-28 03:21:52 +03:00
func WithInsecure() DialOption {
return func(o *dialOptions) {
o.insecure = true
}
}
// WithTransportCredentials returns a DialOption which configures a
// connection level security credentials (e.g., TLS/SSL).
func WithTransportCredentials(creds credentials.TransportCredentials) DialOption {
return func(o *dialOptions) {
o.copts.TransportCredentials = creds
2015-02-06 04:14:05 +03:00
}
}
// WithPerRPCCredentials returns a DialOption which sets
// credentials which will place auth state on each outbound RPC.
func WithPerRPCCredentials(creds credentials.PerRPCCredentials) DialOption {
return func(o *dialOptions) {
o.copts.PerRPCCredentials = append(o.copts.PerRPCCredentials, creds)
2015-03-04 04:08:39 +03:00
}
}
2016-06-06 22:13:00 +03:00
// WithTimeout returns a DialOption that configures a timeout for dialing a ClientConn
// initially. This is valid if and only if WithBlock() is present.
2015-03-04 04:08:39 +03:00
func WithTimeout(d time.Duration) DialOption {
return func(o *dialOptions) {
2016-06-06 22:08:11 +03:00
o.timeout = d
2015-02-06 04:14:05 +03:00
}
}
2015-04-22 02:48:41 +03:00
// WithDialer returns a DialOption that specifies a function to use for dialing network addresses.
func WithDialer(f func(addr string, timeout time.Duration) (net.Conn, error)) DialOption {
2015-04-17 23:50:18 +03:00
return func(o *dialOptions) {
o.copts.Dialer = f
2015-04-17 23:50:18 +03:00
}
}
// WithUserAgent returns a DialOption that specifies a user agent string for all the RPCs.
func WithUserAgent(s string) DialOption {
return func(o *dialOptions) {
o.copts.UserAgent = s
}
}
2015-02-06 04:14:05 +03:00
// Dial creates a client connection the given target.
func Dial(target string, opts ...DialOption) (*ClientConn, error) {
2015-09-29 20:24:03 +03:00
cc := &ClientConn{
target: target,
2016-05-11 05:29:44 +03:00
conns: make(map[Address]*addrConn),
2015-09-29 20:24:03 +03:00
}
2015-02-06 04:14:05 +03:00
for _, opt := range opts {
2015-09-29 20:24:03 +03:00
opt(&cc.dopts)
2015-02-06 04:14:05 +03:00
}
2016-07-05 21:51:54 +03:00
// Set defaults.
2015-10-08 21:05:59 +03:00
if cc.dopts.codec == nil {
cc.dopts.codec = protoCodec{}
}
if cc.dopts.bs == nil {
cc.dopts.bs = DefaultBackoffConfig
}
2016-07-05 21:51:54 +03:00
if cc.dopts.balancer == nil {
cc.dopts.balancer = RoundRobin(nil)
2015-09-29 20:24:03 +03:00
}
2016-07-05 21:51:54 +03:00
if err := cc.dopts.balancer.Start(target); err != nil {
return nil, err
}
2016-06-06 22:08:11 +03:00
var (
ok bool
addrs []Address
)
2016-07-05 21:51:54 +03:00
ch := cc.dopts.balancer.Notify()
if ch == nil {
// There is no name resolver installed.
2016-06-06 22:08:11 +03:00
addrs = append(addrs, Address{Addr: target})
2016-05-07 01:47:09 +03:00
} else {
2016-06-06 22:08:11 +03:00
addrs, ok = <-ch
if !ok || len(addrs) == 0 {
2016-06-06 22:16:33 +03:00
return nil, errNoAddr
2016-05-07 01:47:09 +03:00
}
2016-06-06 22:08:11 +03:00
}
waitC := make(chan error, 1)
2016-06-06 22:08:11 +03:00
go func() {
for _, a := range addrs {
if err := cc.newAddrConn(a, false); err != nil {
2016-06-06 22:08:11 +03:00
waitC <- err
return
2016-05-07 01:47:09 +03:00
}
}
2016-06-06 22:08:11 +03:00
close(waitC)
}()
var timeoutCh <-chan time.Time
if cc.dopts.timeout > 0 {
timeoutCh = time.After(cc.dopts.timeout)
}
select {
case err := <-waitC:
if err != nil {
cc.Close()
return nil, err
}
case <-timeoutCh:
cc.Close()
return nil, ErrClientConnTimeout
}
if ok {
2016-05-26 04:17:23 +03:00
go cc.lbWatcher()
}
2015-10-08 21:05:59 +03:00
colonPos := strings.LastIndex(target, ":")
if colonPos == -1 {
colonPos = len(target)
}
cc.authority = target[:colonPos]
2015-09-29 20:24:03 +03:00
return cc, nil
}
2015-07-31 01:30:26 +03:00
// ConnectivityState indicates the state of a client connection.
type ConnectivityState int
const (
// Idle indicates the ClientConn is idle.
Idle ConnectivityState = iota
// Connecting indicates the ClienConn is connecting.
Connecting
// Ready indicates the ClientConn is ready for work.
Ready
// TransientFailure indicates the ClientConn has seen a failure but expects to recover.
TransientFailure
2015-09-07 07:26:27 +03:00
// Shutdown indicates the ClientConn has started shutting down.
2015-07-31 01:30:26 +03:00
Shutdown
)
2015-08-01 05:00:43 +03:00
func (s ConnectivityState) String() string {
switch s {
case Idle:
return "IDLE"
case Connecting:
return "CONNECTING"
case Ready:
return "READY"
case TransientFailure:
return "TRANSIENT_FAILURE"
case Shutdown:
return "SHUTDOWN"
default:
panic(fmt.Sprintf("unknown connectivity state: %d", s))
}
}
// ClientConn represents a client connection to an RPC server.
2015-02-06 04:14:05 +03:00
type ClientConn struct {
2015-10-08 21:05:59 +03:00
target string
authority string
dopts dialOptions
2016-05-07 01:47:09 +03:00
mu sync.RWMutex
2016-05-11 05:29:44 +03:00
conns map[Address]*addrConn
}
2016-05-26 04:17:23 +03:00
func (cc *ClientConn) lbWatcher() {
2016-07-05 21:51:54 +03:00
for addrs := range cc.dopts.balancer.Notify() {
var (
add []Address // Addresses need to setup connections.
del []*addrConn // Connections need to tear down.
)
cc.mu.Lock()
for _, a := range addrs {
if _, ok := cc.conns[a]; !ok {
add = append(add, a)
2016-05-07 01:47:09 +03:00
}
}
for k, c := range cc.conns {
var keep bool
for _, a := range addrs {
if k == a {
keep = true
break
}
2016-05-07 01:47:09 +03:00
}
if !keep {
del = append(del, c)
2016-05-07 01:47:09 +03:00
}
}
cc.mu.Unlock()
for _, a := range add {
2016-05-26 04:17:23 +03:00
cc.newAddrConn(a, true)
}
for _, c := range del {
2016-05-25 21:28:45 +03:00
c.tearDown(errConnDrain)
2016-05-07 01:47:09 +03:00
}
}
2015-02-06 04:14:05 +03:00
}
func (cc *ClientConn) newAddrConn(addr Address, skipWait bool) error {
2016-05-18 03:18:54 +03:00
ac := &addrConn{
2016-05-11 05:29:44 +03:00
cc: cc,
addr: addr,
dopts: cc.dopts,
shutdownChan: make(chan struct{}),
}
if EnableTracing {
2016-05-18 03:18:54 +03:00
ac.events = trace.NewEventLog("grpc.ClientConn", ac.addr.Addr)
}
2016-05-18 03:18:54 +03:00
if !ac.dopts.insecure {
if ac.dopts.copts.TransportCredentials == nil {
2016-05-25 21:28:45 +03:00
return errNoTransportSecurity
}
} else {
if ac.dopts.copts.TransportCredentials != nil {
return errCredentialsConflict
}
for _, cd := range ac.dopts.copts.PerRPCCredentials {
if cd.RequireTransportSecurity() {
return errTransportCredentialsMissing
}
}
}
2016-05-19 02:26:12 +03:00
// Insert ac into ac.cc.conns. This needs to be done before any getTransport(...) is called.
2016-05-18 03:18:54 +03:00
ac.cc.mu.Lock()
if ac.cc.conns == nil {
ac.cc.mu.Unlock()
return ErrClientConnClosing
}
stale := ac.cc.conns[ac.addr]
2016-05-18 03:18:54 +03:00
ac.cc.conns[ac.addr] = ac
ac.cc.mu.Unlock()
if stale != nil {
// There is an addrConn alive on ac.addr already. This could be due to
// i) stale's Close is undergoing;
// ii) a buggy Balancer notifies duplicated Addresses.
stale.tearDown(errConnDrain)
}
2016-05-18 03:18:54 +03:00
ac.stateCV = sync.NewCond(&ac.mu)
// skipWait may overwrite the decision in ac.dopts.block.
if ac.dopts.block && !skipWait {
2016-05-18 03:18:54 +03:00
if err := ac.resetTransport(false); err != nil {
ac.tearDown(err)
return err
}
// Start to monitor the error status of transport.
2016-05-18 03:18:54 +03:00
go ac.transportMonitor()
} else {
// Start a goroutine connecting to the server asynchronously.
go func() {
2016-05-18 03:18:54 +03:00
if err := ac.resetTransport(false); err != nil {
grpclog.Printf("Failed to dial %s: %v; please retry.", ac.addr.Addr, err)
ac.tearDown(err)
return
}
2016-05-18 03:18:54 +03:00
ac.transportMonitor()
}()
}
return nil
}
2016-05-25 03:19:44 +03:00
func (cc *ClientConn) getTransport(ctx context.Context, opts BalancerGetOptions) (transport.ClientTransport, func(), error) {
2016-07-05 21:51:54 +03:00
addr, put, err := cc.dopts.balancer.Get(ctx, opts)
2016-05-07 01:47:09 +03:00
if err != nil {
return nil, nil, toRPCErr(err)
}
2016-05-07 01:47:09 +03:00
cc.mu.RLock()
2016-05-11 05:29:44 +03:00
if cc.conns == nil {
2016-05-07 01:47:09 +03:00
cc.mu.RUnlock()
2016-06-30 01:21:44 +03:00
return nil, nil, toRPCErr(ErrClientConnClosing)
2016-05-07 01:47:09 +03:00
}
2016-05-11 05:29:44 +03:00
ac, ok := cc.conns[addr]
2016-05-07 01:47:09 +03:00
cc.mu.RUnlock()
if !ok {
if put != nil {
put()
}
return nil, nil, Errorf(codes.Internal, "grpc: failed to find the transport to send the rpc")
2016-05-07 01:47:09 +03:00
}
t, err := ac.wait(ctx, !opts.BlockingWait)
2016-05-07 01:47:09 +03:00
if err != nil {
if put != nil {
put()
}
2016-05-07 01:47:09 +03:00
return nil, nil, err
}
2016-05-07 01:47:09 +03:00
return t, put, nil
}
2016-05-18 03:18:54 +03:00
// Close tears down the ClientConn and all underlying connections.
2016-05-07 01:47:09 +03:00
func (cc *ClientConn) Close() error {
2015-08-01 00:16:02 +03:00
cc.mu.Lock()
2016-05-11 05:29:44 +03:00
if cc.conns == nil {
2016-05-07 01:47:09 +03:00
cc.mu.Unlock()
return ErrClientConnClosing
}
2016-05-11 05:29:44 +03:00
conns := cc.conns
cc.conns = nil
2016-05-07 01:47:09 +03:00
cc.mu.Unlock()
2016-07-05 21:51:54 +03:00
cc.dopts.balancer.Close()
2016-05-11 05:29:44 +03:00
for _, ac := range conns {
ac.tearDown(ErrClientConnClosing)
2016-05-07 01:47:09 +03:00
}
return nil
}
// addrConn is a network connection to a given address.
type addrConn struct {
2016-05-11 05:29:44 +03:00
cc *ClientConn
addr Address
dopts dialOptions
2016-05-07 01:47:09 +03:00
shutdownChan chan struct{}
events trace.EventLog
mu sync.Mutex
state ConnectivityState
stateCV *sync.Cond
down func(error) // the handler called when a connection is down.
// ready is closed and becomes nil when a new transport is up or failed
// due to timeout.
ready chan struct{}
transport transport.ClientTransport
}
// printf records an event in ac's event log, unless ac has been closed.
// REQUIRES ac.mu is held.
func (ac *addrConn) printf(format string, a ...interface{}) {
if ac.events != nil {
ac.events.Printf(format, a...)
}
}
// errorf records an error in ac's event log, unless ac has been closed.
// REQUIRES ac.mu is held.
func (ac *addrConn) errorf(format string, a ...interface{}) {
if ac.events != nil {
ac.events.Errorf(format, a...)
}
}
// getState returns the connectivity state of the Conn
func (ac *addrConn) getState() ConnectivityState {
ac.mu.Lock()
defer ac.mu.Unlock()
return ac.state
}
// waitForStateChange blocks until the state changes to something other than the sourceState.
func (ac *addrConn) waitForStateChange(ctx context.Context, sourceState ConnectivityState) (ConnectivityState, error) {
ac.mu.Lock()
defer ac.mu.Unlock()
if sourceState != ac.state {
return ac.state, nil
2015-08-03 21:29:27 +03:00
}
2015-08-03 21:45:42 +03:00
done := make(chan struct{})
var err error
2015-08-03 23:18:25 +03:00
go func() {
2015-08-01 00:16:02 +03:00
select {
case <-ctx.Done():
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
err = ctx.Err()
2016-05-07 01:47:09 +03:00
ac.stateCV.Broadcast()
ac.mu.Unlock()
2015-08-01 00:16:02 +03:00
case <-done:
}
2015-08-03 23:18:25 +03:00
}()
2015-08-01 05:00:43 +03:00
defer close(done)
2016-05-07 01:47:09 +03:00
for sourceState == ac.state {
ac.stateCV.Wait()
if err != nil {
2016-05-07 01:47:09 +03:00
return ac.state, err
2015-08-01 05:00:43 +03:00
}
2015-08-01 00:16:02 +03:00
}
2016-05-07 01:47:09 +03:00
return ac.state, nil
2015-08-01 00:16:02 +03:00
}
2016-05-07 01:47:09 +03:00
func (ac *addrConn) resetTransport(closeTransport bool) error {
2015-02-06 04:14:05 +03:00
var retries int
for {
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
ac.printf("connecting")
if ac.state == Shutdown {
// ac.tearDown(...) has been invoked.
ac.mu.Unlock()
2016-05-25 21:28:45 +03:00
return errConnClosing
2015-02-06 04:14:05 +03:00
}
2016-05-07 01:47:09 +03:00
if ac.down != nil {
2016-05-25 21:28:45 +03:00
ac.down(downErrorf(false, true, "%v", errNetworkIO))
2016-05-07 01:47:09 +03:00
ac.down = nil
}
ac.state = Connecting
ac.stateCV.Broadcast()
t := ac.transport
ac.mu.Unlock()
if closeTransport && t != nil {
t.Close()
2015-02-06 04:14:05 +03:00
}
2016-05-07 01:47:09 +03:00
sleepTime := ac.dopts.bs.backoff(retries)
2016-06-06 22:08:11 +03:00
ac.dopts.copts.Timeout = sleepTime
if sleepTime < minConnectTimeout {
ac.dopts.copts.Timeout = minConnectTimeout
2015-07-28 21:12:07 +03:00
}
connectTime := time.Now()
2016-06-06 22:08:11 +03:00
newTransport, err := transport.NewClientTransport(ac.addr.Addr, &ac.dopts.copts)
2015-02-06 04:14:05 +03:00
if err != nil {
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
if ac.state == Shutdown {
// ac.tearDown(...) has been invoked.
ac.mu.Unlock()
2016-05-25 21:28:45 +03:00
return errConnClosing
}
2016-05-07 01:47:09 +03:00
ac.errorf("transient failure: %v", err)
ac.state = TransientFailure
ac.stateCV.Broadcast()
if ac.ready != nil {
close(ac.ready)
ac.ready = nil
2015-09-25 23:21:25 +03:00
}
2016-05-07 01:47:09 +03:00
ac.mu.Unlock()
2015-07-28 21:12:07 +03:00
sleepTime -= time.Since(connectTime)
if sleepTime < 0 {
sleepTime = 0
}
2015-02-06 04:14:05 +03:00
closeTransport = false
select {
case <-time.After(sleepTime):
2016-05-13 04:52:24 +03:00
case <-ac.shutdownChan:
}
2015-02-06 04:14:05 +03:00
retries++
2016-05-07 01:47:09 +03:00
grpclog.Printf("grpc: addrConn.resetTransport failed to create client transport: %v; Reconnecting to %q", err, ac.addr)
2015-02-06 04:14:05 +03:00
continue
}
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
ac.printf("ready")
if ac.state == Shutdown {
// ac.tearDown(...) has been invoked.
ac.mu.Unlock()
2015-03-05 20:45:50 +03:00
newTransport.Close()
2016-05-25 21:28:45 +03:00
return errConnClosing
2015-03-05 20:45:50 +03:00
}
2016-05-07 01:47:09 +03:00
ac.state = Ready
ac.stateCV.Broadcast()
ac.transport = newTransport
if ac.ready != nil {
close(ac.ready)
ac.ready = nil
2015-02-06 04:14:05 +03:00
}
2016-07-05 21:51:54 +03:00
ac.down = ac.cc.dopts.balancer.Up(ac.addr)
2016-05-07 01:47:09 +03:00
ac.mu.Unlock()
2015-02-06 04:14:05 +03:00
return nil
}
}
// Run in a goroutine to track the error in transport and create the
// new transport if an error happens. It returns when the channel is closing.
2016-05-07 01:47:09 +03:00
func (ac *addrConn) transportMonitor() {
2015-02-06 04:14:05 +03:00
for {
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
t := ac.transport
ac.mu.Unlock()
2015-02-06 04:14:05 +03:00
select {
2015-07-31 01:30:26 +03:00
// shutdownChan is needed to detect the teardown when
2016-05-07 01:47:09 +03:00
// the addrConn is idle (i.e., no RPC in flight).
case <-ac.shutdownChan:
2015-02-06 04:14:05 +03:00
return
2016-07-22 02:19:34 +03:00
case <-t.GoAway():
ac.tearDown(errConnDrain)
ac.cc.newAddrConn(ac.addr, true)
return
case <-t.Error():
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
if ac.state == Shutdown {
2016-07-28 03:27:10 +03:00
// ac has been shutdown.
2016-05-07 01:47:09 +03:00
ac.mu.Unlock()
return
}
2016-05-07 01:47:09 +03:00
ac.state = TransientFailure
ac.stateCV.Broadcast()
ac.mu.Unlock()
if err := ac.resetTransport(true); err != nil {
ac.mu.Lock()
ac.printf("transport exiting: %v", err)
ac.mu.Unlock()
grpclog.Printf("grpc: addrConn.transportMonitor exits due to: %v", err)
2015-02-06 04:14:05 +03:00
return
}
}
}
}
// wait blocks until i) the new transport is up or ii) ctx is done or iii) ac is closed or
// iv) transport is in TransientFailure and the RPC is fail-fast.
func (ac *addrConn) wait(ctx context.Context, failFast bool) (transport.ClientTransport, error) {
2015-02-06 04:14:05 +03:00
for {
2016-05-07 01:47:09 +03:00
ac.mu.Lock()
2015-02-06 04:14:05 +03:00
switch {
2016-05-07 01:47:09 +03:00
case ac.state == Shutdown:
ac.mu.Unlock()
2016-05-25 21:28:45 +03:00
return nil, errConnClosing
2016-05-07 01:47:09 +03:00
case ac.state == Ready:
ct := ac.transport
ac.mu.Unlock()
return ct, nil
case ac.state == TransientFailure && failFast:
ac.mu.Unlock()
return nil, Errorf(codes.Unavailable, "grpc: RPC failed fast due to transport failure")
2015-02-06 04:14:05 +03:00
default:
2016-05-07 01:47:09 +03:00
ready := ac.ready
2015-02-06 04:14:05 +03:00
if ready == nil {
ready = make(chan struct{})
2016-05-07 01:47:09 +03:00
ac.ready = ready
2015-02-06 04:14:05 +03:00
}
2016-05-07 01:47:09 +03:00
ac.mu.Unlock()
2015-02-06 04:14:05 +03:00
select {
case <-ctx.Done():
return nil, toRPCErr(ctx.Err())
// Wait until the new transport is ready or failed.
2015-02-06 04:14:05 +03:00
case <-ready:
}
}
}
}
// tearDown starts to tear down the addrConn.
2015-02-06 04:14:05 +03:00
// TODO(zhaoq): Make this synchronous to avoid unbounded memory consumption in
2016-05-07 01:47:09 +03:00
// some edge cases (e.g., the caller opens and closes many addrConn's in a
2015-02-06 04:14:05 +03:00
// tight loop.
2016-05-07 01:47:09 +03:00
func (ac *addrConn) tearDown(err error) {
ac.mu.Lock()
defer func() {
ac.mu.Unlock()
ac.cc.mu.Lock()
if ac.cc.conns != nil {
delete(ac.cc.conns, ac.addr)
}
ac.cc.mu.Unlock()
}()
2016-05-27 00:53:32 +03:00
if ac.down != nil {
ac.down(downErrorf(false, false, "%v", err))
ac.down = nil
}
if err == errConnDrain && ac.transport != nil {
// GracefulClose(...) may be executed multiple times when
// i) receiving multiple GoAway frames from the server; or
// ii) there are concurrent name resolver/Balancer triggered
// address removal and GoAway.
ac.transport.GracefulClose()
}
if ac.state == Shutdown {
return
}
ac.state = Shutdown
2016-05-07 01:47:09 +03:00
ac.stateCV.Broadcast()
if ac.events != nil {
ac.events.Finish()
ac.events = nil
2015-03-05 00:00:47 +03:00
}
2016-05-07 01:47:09 +03:00
if ac.ready != nil {
close(ac.ready)
ac.ready = nil
}
if ac.transport != nil && err != errConnDrain {
ac.transport.Close()
}
2016-05-07 01:47:09 +03:00
if ac.shutdownChan != nil {
close(ac.shutdownChan)
}
return
2015-02-06 04:14:05 +03:00
}