mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-09 22:51:02 +00:00
reset grpc connection after TransientFailure
This commit is contained in:
parent
7aef0743e1
commit
5ca37bf3f3
1 changed files with 74 additions and 12 deletions
|
|
@ -17,53 +17,113 @@ package client
|
|||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"sync"
|
||||
|
||||
pb "github.com/pilosa/pilosa/v2/proto"
|
||||
"github.com/pkg/errors"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/connectivity"
|
||||
"google.golang.org/grpc/credentials"
|
||||
)
|
||||
|
||||
// GRPCClient is a client for working with the gRPC server.
|
||||
type GRPCClient struct {
|
||||
dialTarget string
|
||||
tlsConfig *tls.Config
|
||||
|
||||
mu sync.RWMutex
|
||||
conn *grpc.ClientConn
|
||||
}
|
||||
|
||||
// NewGRPCClient returns a new instance of GRPCClient.
|
||||
func NewGRPCClient(dialTarget string, tlsConfig *tls.Config) (*GRPCClient, error) {
|
||||
c := &GRPCClient{
|
||||
dialTarget: dialTarget,
|
||||
tlsConfig: tlsConfig,
|
||||
}
|
||||
// resetConn sets GRPCClient.conn when it doesn't
|
||||
// exist yet.
|
||||
if err := c.resetConn(); err != nil {
|
||||
return nil, errors.Wrap(err, "setting connection")
|
||||
}
|
||||
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// resetConn resets the gRPC client connection. This method
|
||||
// can also be used to initially set the client connection
|
||||
// because it only tries to first close the connection if
|
||||
// the connection already exists.
|
||||
func (c *GRPCClient) resetConn() error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
// If an existing connection exists, close it first.
|
||||
if c.conn != nil {
|
||||
if err := c.conn.Close(); err != nil {
|
||||
return errors.Wrap(err, "closing existing connection")
|
||||
}
|
||||
}
|
||||
|
||||
var opts []grpc.DialOption
|
||||
if tlsConfig != nil {
|
||||
creds := credentials.NewTLS(tlsConfig)
|
||||
if c.tlsConfig != nil {
|
||||
creds := credentials.NewTLS(c.tlsConfig)
|
||||
opts = append(opts, grpc.WithTransportCredentials(creds))
|
||||
} else {
|
||||
opts = append(opts, grpc.WithInsecure())
|
||||
}
|
||||
|
||||
gconn, err := grpc.Dial(dialTarget, opts...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "creating new grpc client")
|
||||
var err error
|
||||
if c.conn, err = grpc.Dial(c.dialTarget, opts...); err != nil {
|
||||
return errors.Wrap(err, "creating new grpc client")
|
||||
}
|
||||
|
||||
return &GRPCClient{
|
||||
conn: gconn,
|
||||
}, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close closes any connections the client has opened.
|
||||
func (c *GRPCClient) Close() error {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
|
||||
if c.conn != nil {
|
||||
return c.conn.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Conn returns the gRPC client connection. If the connection
|
||||
// has gone into state `TransientFailure`, this method tries
|
||||
// to reset the connection and return that new connection.
|
||||
func (c *GRPCClient) Conn() *grpc.ClientConn {
|
||||
c.mu.RLock()
|
||||
if c.conn == nil {
|
||||
c.mu.RUnlock()
|
||||
return nil
|
||||
} else if c.conn.GetState() != connectivity.TransientFailure {
|
||||
defer c.mu.RUnlock()
|
||||
return c.conn
|
||||
}
|
||||
c.mu.RUnlock()
|
||||
|
||||
if err := c.resetConn(); err != nil {
|
||||
// TODO: log this error once logger is in place
|
||||
}
|
||||
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
return c.conn
|
||||
}
|
||||
|
||||
// Query returns a stream of RowResponse for the given index and PQL string.
|
||||
func (c *GRPCClient) Query(ctx context.Context, index string, pql string) (pb.StreamClient, error) {
|
||||
if c.conn == nil {
|
||||
conn := c.Conn()
|
||||
|
||||
if conn == nil {
|
||||
return nil, errors.New("client has not established a grpc connection")
|
||||
}
|
||||
|
||||
grpcClient := pb.NewPilosaClient(c.conn)
|
||||
grpcClient := pb.NewPilosaClient(conn)
|
||||
|
||||
stream, err := grpcClient.QueryPQL(ctx, &pb.QueryPQLRequest{
|
||||
Index: index,
|
||||
|
|
@ -81,7 +141,9 @@ func (c *GRPCClient) Query(ctx context.Context, index string, pql string) (pb.St
|
|||
// Inspect returns a stream of RowResponse for the given index, columns, and filters.
|
||||
// It is intended to mimic something like "select [fields] from table where recordID IN (...)".
|
||||
func (c *GRPCClient) Inspect(ctx context.Context, index string, columnIDs []uint64, columnKeys []string, fieldFilters []string, limit, offset uint64) (pb.StreamClient, error) {
|
||||
if c.conn == nil {
|
||||
conn := c.Conn()
|
||||
|
||||
if conn == nil {
|
||||
return nil, errors.New("client has not established a grpc connection")
|
||||
}
|
||||
|
||||
|
|
@ -97,7 +159,7 @@ func (c *GRPCClient) Inspect(ctx context.Context, index string, columnIDs []uint
|
|||
idsOrKeys.Type = &pb.IdsOrKeys_Ids{Ids: &pb.Uint64Array{Vals: columnIDs}}
|
||||
}
|
||||
|
||||
grpcClient := pb.NewPilosaClient(c.conn)
|
||||
grpcClient := pb.NewPilosaClient(conn)
|
||||
|
||||
stream, err := grpcClient.Inspect(ctx, &pb.InspectRequest{
|
||||
Index: index,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue