Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
134 changes: 33 additions & 101 deletions bigtable/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ import (
"google.golang.org/grpc/metadata"
)

const directpathEnvVar = "CBT_ENABLE_DIRECTPATH"
const directAccessEnvVar = "CBT_ENABLE_DIRECTPATH"

// Client is a client for reading and writing data to tables in an instance.
//
Expand All @@ -50,8 +50,7 @@ type Client struct {
retryOption gax.CallOption
executeQueryRetryOption gax.CallOption
featureFlagsMD metadata.MD // Pre-computed feature flags metadata to be sent with each request.
dynamicScaleMonitor *btransport.DynamicScaleMonitor
connsRecycler *btransport.ConnectionRecycler
mPool btransport.ManagedChannelPool
}

// ClientConfig has configurations for the client.
Expand Down Expand Up @@ -131,7 +130,7 @@ func NewClientWithConfig(ctx context.Context, project, instance string, config C
option.WithGRPCDialOption(grpc.WithDefaultCallOptions(grpc.MaxCallSendMsgSize(1<<28), grpc.MaxCallRecvMsgSize(1<<28))),
)

var directPathOptions = []option.ClientOption{
var directAccessOptions = []option.ClientOption{
internaloption.EnableDirectPath(true),
internaloption.EnableDirectPathXds(),
internaloption.AllowHardBoundTokens("ALTS"),
Expand Down Expand Up @@ -165,10 +164,7 @@ func NewClientWithConfig(ctx context.Context, project, instance string, config C
allowDirectAccess := isDirectAccessEnabled(config)
directAccessMD := createFeatureFlagsMD(metricsTracerFactory.enabled, disableRetryInfo, allowDirectAccess)

var connPool gtransport.ConnPool
var dsm *btransport.DynamicScaleMonitor
var connRecycler *btransport.ConnectionRecycler

var mPool btransport.ManagedChannelPool
enableBigtableConnPool := btopt.EnableBigtableConnectionPool()
grpcConnOptType := reflect.TypeOf(option.WithGRPCConn(nil))
for _, opt := range opts {
Expand All @@ -177,97 +173,40 @@ func NewClientWithConfig(ctx context.Context, project, instance string, config C
break
}
}
var connPoolSize int
if !enableBigtableConnPool {
// Use the regular ConnPool
// For regular ConnPool the Direct Access is off by default so we need to check the env var again.
if enabled, _ := strconv.ParseBool(os.Getenv(directpathEnvVar)); enabled {
o = append(o, directPathOptions...)
}
regConnPool, err := gtransport.DialPool(ctx, o...)
if err != nil {
return nil, err
}
connPool = regConnPool
} else { // Use the BigtableConnPool
uResolver, err := internaloption.NewUnsafeResolver(o...)
if err != nil {
// just fallback
connPoolSize = defaultBigtableConnPoolSize
}

connPoolSize = uResolver.ResolvedGRPCConnPoolSize()
// Fallback to 10 if it resolves to 0
if connPoolSize == 0 {
connPoolSize = defaultBigtableConnPoolSize
if enabled, _ := strconv.ParseBool(os.Getenv(directAccessEnvVar)); enabled {
o = append(o, directAccessOptions...)
}
}

fullInstanceName := fmt.Sprintf("projects/%s/instances/%s", project, instance)

var poolOpts []btransport.BigtableChannelPoolOption
poolOpts = append(poolOpts,
btransport.WithInstanceName(fullInstanceName),
btransport.WithAppProfile(config.AppProfile),
btransport.WithFeatureFlagsMetadata(directAccessMD),
btransport.WithMetricsReporterConfig(btopt.DefaultMetricsReporterConfig()),
btransport.WithMeterProvider(metricsTracerFactory.otelMeterProvider),
btransport.WithDirectAccessFeatureFlagsMetadata(directAccessMD),
)

// Only setup DirectPath dialers if not disabled by config/env
if allowDirectAccess {
directAccessDialerOptions := make([]option.ClientOption, len(o))
copy(directAccessDialerOptions, o)
directAccessDialerOptions = append(directAccessDialerOptions, directPathOptions...)
directAccessDialer := func() (*btransport.BigtableConn, error) {
grpcConn, err := gtransport.Dial(ctx, directAccessDialerOptions...)
if err != nil {
return nil, err
}
return btransport.NewBigtableConn(grpcConn), nil
}
poolOpts = append(poolOpts, btransport.WithDirectAccessDialer(directAccessDialer))
}

btPool, err := btransport.NewBigtableChannelPool(ctx,
connPoolSize,
btopt.BigtableLoadBalancingStrategy(),
func() (*btransport.BigtableConn, error) {
grpcConn, err := gtransport.Dial(ctx, o...)
if err != nil {
return nil, err
}
return btransport.NewBigtableConn(grpcConn), nil
},
clientCreationTimestamp,
// options
poolOpts...,
)
if err != nil {
return nil, err
}

connPool = btPool

// Validate dynamic config early if enabled
if !config.DisableDynamicChannelPool {
if err := btransport.ValidateDynamicConfig(btopt.DefaultDynamicChannelPoolConfig(), defaultBigtableConnPoolSize); err != nil {
return nil, fmt.Errorf("invalid DynamicChannelPoolConfig: %w", err)
}
poolConfig := btransport.ChannelPoolConfig{
AppProfile: config.AppProfile,
DisableDynamicChannelPool: config.DisableDynamicChannelPool,
DisableConnectionRecycler: config.DisableConnectionRecycler,
DisableDirectAccess: config.DisableDirectAccess,
}

dsm = btransport.NewDynamicScaleMonitor(btopt.DefaultDynamicChannelPoolConfig(), btPool)
dsm.Start(ctx) // Start the monitor's background goroutine
}
// connection recyler.
if !config.DisableConnectionRecycler {
connRecycler = btransport.NewConnectionRecycler(btopt.DefaultConnectionRecycleConfig(), btPool)
connRecycler.Start(ctx) // Start the monitor's background goroutine
}
mPool, err = btransport.CreateAndStartManagedChannelPool(
ctx,
project,
instance,
poolConfig,
metricsTracerFactory.otelMeterProvider,
o,
directAccessOptions,
directAccessMD,
clientCreationTimestamp,
enableBigtableConnPool,
)
if err != nil {
return nil, err
}

return &Client{
connPool: connPool,
client: btpb.NewBigtableClient(connPool),
connPool: mPool.Pool,
client: btpb.NewBigtableClient(mPool.Pool),
project: project,
instance: instance,
appProfile: config.AppProfile,
Expand All @@ -276,23 +215,16 @@ func NewClientWithConfig(ctx context.Context, project, instance string, config C
retryOption: retryOption,
executeQueryRetryOption: executeQueryRetryOption,
featureFlagsMD: directAccessMD,
dynamicScaleMonitor: dsm,
connsRecycler: connRecycler,
mPool: mPool,
}, nil
}

// Close closes the Client.
func (c *Client) Close() error {
if c.dynamicScaleMonitor != nil {
c.dynamicScaleMonitor.Stop()
}
if c.metricsTracerFactory != nil {
c.metricsTracerFactory.shutdown()
}
if c.connsRecycler != nil {
c.connsRecycler.Stop()
}
return c.connPool.Close()
return c.mPool.Close()
}

func (c *Client) fullInstanceName() string {
Expand Down Expand Up @@ -408,9 +340,9 @@ func (c *Client) newBuiltinMetricsTracer(ctx context.Context, table string, isSt
}

func isDirectAccessEnabled(config ClientConfig) bool {
if os.Getenv(directpathEnvVar) == "" {
if os.Getenv(directAccessEnvVar) == "" {
return !config.DisableDirectAccess
}
res, _ := strconv.ParseBool(os.Getenv(directpathEnvVar))
res, _ := strconv.ParseBool(os.Getenv(directAccessEnvVar))
return res
}
Loading
Loading