Skip to content

Commit 5214ab7

Browse files
authored
feat(bigtable): modularize channel priming behind a ChannelPrimer interface (#20027)
## Summary Mirrors the pluggable shape that PR #19987 established for the Direct Access compatibility check. Channel priming had been hard-wired into the pool via three loose options (`WithInstanceName` / `WithAppProfile` / `WithFeatureFlagsMetadata`) that the `connectionFactory` had to stitch back together on every dial. - New `ChannelPrimer` interface (`channel_primer.go`): `Prime(ctx, *BigtableConn) error`. The pluggable extension point that future pool factories (session-based, custom) can swap in. - `pingAndWarmChannelPrimer` is today's only implementation; it owns the `(instance, appProfile, featureFlagsMD)` tuple and delegates to `BigtableConn.Prime` — the existing `PingAndWarm` RPC stays untouched. - New `WithChannelPrimer` pool option replaces `WithInstanceName` / `WithAppProfile` / `WithFeatureFlagsMetadata`. - `connectionFactory` now carries a `ChannelPrimer` instead of three individual fields. **When the primer is nil, `primeWithRetry` returns immediately** — the pool dials the channel and puts it straight into rotation, no `PingAndWarm` sent. The classic channel pool factory still wires the `pingAndWarmChannelPrimer`, so user-facing default behavior is unchanged. - `pingAndWarmDirectAccessChecker` reuses the same primer rather than duplicating `conn.Prime(ctx, instance, profile, flags)` in two places. The factory now constructs **one** primer and shares it with both the pool (via `WithChannelPrimer`) and the checker, eliminating the three-arg drift that existed across the two consumers. ## Test plan - [x] `go build ./...` clean - [x] `go vet ./...` clean - [x] `golint` clean on touched files - [x] `go test ./internal/transport/...` passes (existing tests migrated; new `channel_primer_test.go` covers (a) `pingAndWarmChannelPrimer.Prime` issues `PingAndWarm` carrying the configured feature-flag metadata, and (b) `connectionFactory` skips priming entirely when the primer is nil)
1 parent 2e6820b commit 5214ab7

6 files changed

Lines changed: 214 additions & 76 deletions

File tree

‎bigtable/internal/transport/channel_pool_factory.go‎

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -153,15 +153,17 @@ func CreateBigtableChannelPool(
153153

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

156-
// directAccessMD is the feature-flag metadata used for priming on both
157-
// the direct-access and standard-path connection factories — the pool
158-
// holds it via WithFeatureFlagsMetadata and applies it to every Prime().
156+
// Build a single PingAndWarm primer and share it with both the pool's
157+
// connection factory (via WithChannelPrimer) and the direct-access
158+
// compatibility checker. Keeping the (instanceName, appProfile,
159+
// featureFlagsMD) tuple in one place avoids the three-arg drift between
160+
// the two consumers.
161+
primer := newPingAndWarmChannelPrimer(fullInstanceName, config.AppProfile, directAccessMD)
162+
159163
poolOpts := []BigtableChannelPoolOption{
160-
WithInstanceName(fullInstanceName),
161-
WithAppProfile(config.AppProfile),
162164
WithMetricsReporterConfig(btopt.DefaultMetricsReporterConfig()),
163165
WithMeterProvider(otelMeterProvider),
164-
WithFeatureFlagsMetadata(directAccessMD),
166+
WithChannelPrimer(primer),
165167
}
166168

167169
// Pluggable Direct Access strategy: the classic channel pool factory uses
@@ -184,9 +186,7 @@ func CreateBigtableChannelPool(
184186
}
185187
checker := newPingAndWarmDirectAccessChecker(
186188
directAccessDialer,
187-
fullInstanceName,
188-
config.AppProfile,
189-
directAccessMD,
189+
primer,
190190
otelMeterProvider,
191191
nil, // logger plumbed by callers once available
192192
)
Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,63 @@
1+
// Copyright 2026 Google LLC
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package internal
16+
17+
import (
18+
"context"
19+
20+
"google.golang.org/grpc/metadata"
21+
)
22+
23+
// ChannelPrimer warms a freshly-dialed Bigtable channel before it is put
24+
// into rotation. The pool's connection factory consults the registered
25+
// primer after every successful dial; a nil ChannelPrimer means the pool
26+
// skips priming entirely and hands the raw connection straight to the
27+
// pool.
28+
//
29+
// Implementations differ in HOW the channel is warmed:
30+
// - pingAndWarmChannelPrimer issues a PingAndWarm against the configured
31+
// instance / app profile with the supplied feature-flag metadata
32+
// (today's only behavior, used by the classic channel pool factory).
33+
type ChannelPrimer interface {
34+
// Prime warms conn so the next request served by it does not pay the
35+
// first-RPC connection-setup cost. The factory wraps Prime in a retry
36+
// loop, so transient errors should propagate as-is.
37+
Prime(ctx context.Context, conn *BigtableConn) error
38+
}
39+
40+
// pingAndWarmChannelPrimer primes a channel by issuing a PingAndWarm RPC
41+
// against the configured instance + app profile, carrying the supplied
42+
// feature-flag metadata. Stateless aside from the configured identifiers,
43+
// so a single primer instance is shared across all dials in a pool.
44+
type pingAndWarmChannelPrimer struct {
45+
instanceName string
46+
appProfile string
47+
featureFlagsMD metadata.MD
48+
}
49+
50+
// newPingAndWarmChannelPrimer constructs the today-default channel primer.
51+
func newPingAndWarmChannelPrimer(instanceName, appProfile string, featureFlagsMD metadata.MD) *pingAndWarmChannelPrimer {
52+
return &pingAndWarmChannelPrimer{
53+
instanceName: instanceName,
54+
appProfile: appProfile,
55+
featureFlagsMD: featureFlagsMD,
56+
}
57+
}
58+
59+
// Prime delegates to BigtableConn.Prime, which sends PingAndWarm and
60+
// records the ALTS / IP-protocol observations on conn as a side effect.
61+
func (p *pingAndWarmChannelPrimer) Prime(ctx context.Context, conn *BigtableConn) error {
62+
return conn.Prime(ctx, p.instanceName, p.appProfile, p.featureFlagsMD)
63+
}
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
// Copyright 2026 Google LLC
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package internal
16+
17+
import (
18+
"context"
19+
"testing"
20+
21+
"google.golang.org/grpc/metadata"
22+
)
23+
24+
// TestPingAndWarmChannelPrimer_Prime verifies the primer delegates to
25+
// BigtableConn.Prime carrying the configured instance name, app profile,
26+
// and feature-flag metadata — i.e. that the primer is a thin wrapper that
27+
// keeps the (instance, profile, flags) tuple in one place.
28+
func TestPingAndWarmChannelPrimer_Prime(t *testing.T) {
29+
fake := &fakeService{}
30+
addr := setupTestServer(t, fake)
31+
conn, err := dialBigtableserver(addr)
32+
if err != nil {
33+
t.Fatalf("dial: %v", err)
34+
}
35+
t.Cleanup(func() { conn.Close() })
36+
37+
flagsMD := metadata.Pairs("bigtable-features", "primer-test")
38+
primer := newPingAndWarmChannelPrimer(testInstanceName, testAppProfile, flagsMD)
39+
40+
if err := primer.Prime(context.Background(), conn); err != nil {
41+
t.Fatalf("Prime returned error: %v", err)
42+
}
43+
44+
if got := fake.getPingCount(); got != 1 {
45+
t.Errorf("PingAndWarm call count = %d, want 1", got)
46+
}
47+
48+
gotMD := fake.getPrimeMetadata()
49+
if got := gotMD.Get("bigtable-features"); len(got) != 1 || got[0] != "primer-test" {
50+
t.Errorf("feature-flag metadata on PingAndWarm = %v, want [primer-test]", got)
51+
}
52+
if got := gotMD.Get("x-goog-request-params"); len(got) != 1 {
53+
t.Errorf("x-goog-request-params on PingAndWarm = %v, want one entry derived from instance/profile", got)
54+
}
55+
}
56+
57+
// TestConnectionFactory_NilPrimerSkipsPriming verifies the contract that a
58+
// nil ChannelPrimer turns priming off: newEntry dials the channel and
59+
// returns it without issuing PingAndWarm.
60+
func TestConnectionFactory_NilPrimerSkipsPriming(t *testing.T) {
61+
fake := &fakeService{}
62+
addr := setupTestServer(t, fake)
63+
64+
factory := &connectionFactory{
65+
dial: func() (*BigtableConn, error) { return dialBigtableserver(addr) },
66+
primer: nil,
67+
}
68+
69+
entry, err := factory.newEntry(context.Background())
70+
if err != nil {
71+
t.Fatalf("newEntry returned error: %v", err)
72+
}
73+
t.Cleanup(func() { entry.conn.Close() })
74+
75+
if got := fake.getPingCount(); got != 0 {
76+
t.Errorf("PingAndWarm call count with nil primer = %d, want 0", got)
77+
}
78+
}

‎bigtable/internal/transport/connpool.go‎

Lines changed: 41 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -101,13 +101,6 @@ type connPoolStats struct {
101101

102102
var _ Monitor = (*MetricsReporter)(nil)
103103

104-
// WithAppProfile provides the appProfile
105-
func WithAppProfile(appProfile string) BigtableChannelPoolOption {
106-
return func(p *BigtableChannelPool) {
107-
p.appProfile = appProfile
108-
}
109-
}
110-
111104
// WithMeterProvider provides the meter provider for writing metrics
112105
func WithMeterProvider(mp metric.MeterProvider) BigtableChannelPoolOption {
113106
return func(p *BigtableChannelPool) {
@@ -129,24 +122,22 @@ func WithDirectAccessChecker(checker DirectAccessChecker) BigtableChannelPoolOpt
129122
}
130123
}
131124

132-
// WithLogger provides the logger for logging events
133-
func WithLogger(logger *log.Logger) BigtableChannelPoolOption {
125+
// WithChannelPrimer plugs in the strategy used to warm freshly-dialed
126+
// channels before they enter rotation. Optional: when no primer is supplied,
127+
// the pool's connection factory dials the channel and returns it without
128+
// issuing any prime RPC. The classic channel pool factory wires up a
129+
// PingAndWarm-based primer; alternative pool factories can swap in a
130+
// different strategy (e.g. session-based) or pass nothing at all.
131+
func WithChannelPrimer(primer ChannelPrimer) BigtableChannelPoolOption {
134132
return func(p *BigtableChannelPool) {
135-
p.logger = logger
133+
p.channelPrimer = primer
136134
}
137135
}
138136

139-
// WithInstanceName provides the full instance Name
140-
func WithInstanceName(instanceName string) BigtableChannelPoolOption {
141-
return func(p *BigtableChannelPool) {
142-
p.instanceName = instanceName
143-
}
144-
}
145-
146-
// WithFeatureFlagsMetadata provides the feature flags metadata
147-
func WithFeatureFlagsMetadata(featureFlagsMd metadata.MD) BigtableChannelPoolOption {
137+
// WithLogger provides the logger for logging events
138+
func WithLogger(logger *log.Logger) BigtableChannelPoolOption {
148139
return func(p *BigtableChannelPool) {
149-
p.featureFlagsMD = featureFlagsMd
140+
p.logger = logger
150141
}
151142
}
152143

@@ -372,10 +363,7 @@ type BigtableChannelPool struct {
372363
poolCtx context.Context // Context for the pool's background tasks
373364
poolCancel context.CancelFunc // Function to cancel the poolCtx
374365

375-
logger *log.Logger // logging events
376-
appProfile string
377-
instanceName string
378-
featureFlagsMD metadata.MD
366+
logger *log.Logger // logging events
379367

380368
factory *connectionFactory // Use the factory for connection creation
381369

@@ -391,6 +379,13 @@ type BigtableChannelPool struct {
391379
// direct_access/compatible metric still surfaces the off state.
392380
directAccessChecker DirectAccessChecker
393381

382+
// channelPrimer is the pluggable strategy used to warm freshly-dialed
383+
// channels. Optional: when nil, the connection factory skips priming
384+
// entirely and hands the raw connection straight to the pool. The
385+
// classic channel pool factory wires up a PingAndWarm-based primer; the
386+
// future session-pool factory may skip it.
387+
channelPrimer ChannelPrimer
388+
394389
// background monitors
395390
monitors []Monitor
396391
}
@@ -441,9 +436,9 @@ func NewBigtableChannelPool(ctx context.Context, connPoolSize int, strategy btop
441436

442437
// Default to the standard dialer. The Direct Access checker may swap the
443438
// dialer for the direct-access equivalent after a successful compatibility
444-
// probe. Feature-flag metadata always comes from pool.featureFlagsMD (set
445-
// via WithFeatureFlagsMetadata) — both the direct-access and standard-path
446-
// connection factories read it from the same place.
439+
// probe. The ChannelPrimer (if any) is the single source of priming
440+
// behavior — both the direct-access and standard-path factories run
441+
// fresh connections through it before they enter rotation.
447442
factoryDial := dial
448443

449444
var firstConn *BigtableConn
@@ -463,11 +458,9 @@ func NewBigtableChannelPool(ctx context.Context, connPoolSize int, strategy btop
463458

464459
// Initialize the connectionFactory
465460
pool.factory = &connectionFactory{
466-
dial: factoryDial,
467-
instanceName: pool.instanceName,
468-
appProfile: pool.appProfile,
469-
featureFlagsMD: pool.featureFlagsMD,
470-
logger: pool.logger,
461+
dial: factoryDial,
462+
primer: pool.channelPrimer,
463+
logger: pool.logger,
471464
}
472465

473466
// Set the selection function based on the strategy
@@ -980,18 +973,18 @@ func (p *BigtableChannelPool) removeConnections(decreaseDelta, minConns, maxRemo
980973

981974
}
982975

983-
// connectionFactory is responsible for creating and priming new Bigtable connections.
984-
// TODO remove these members from BigtableConnPool struct
976+
// connectionFactory is responsible for creating and (optionally) priming
977+
// new Bigtable connections. When primer is nil the factory dials and
978+
// returns the connection without warming it.
985979
type connectionFactory struct {
986-
dial func() (*BigtableConn, error)
987-
instanceName string
988-
appProfile string
989-
featureFlagsMD metadata.MD
990-
logger *log.Logger
980+
dial func() (*BigtableConn, error)
981+
primer ChannelPrimer
982+
logger *log.Logger
991983
}
992984

993-
// newEntry creates a new connection, primes it, and returns it as a connEntry.
994-
// Blocks until the connection is successfully primed, or returns an error.
985+
// newEntry creates a new connection, primes it (if a primer is configured),
986+
// and returns it as a connEntry. Blocks until the connection is ready, or
987+
// returns an error.
995988
func (cf *connectionFactory) newEntry(ctx context.Context) (*connEntry, error) {
996989
conn, err := cf.dial()
997990
if err != nil {
@@ -1006,8 +999,13 @@ func (cf *connectionFactory) newEntry(ctx context.Context) (*connEntry, error) {
1006999
return &connEntry{conn: conn}, nil
10071000
}
10081001

1009-
// primeWithRetry attempts to prime the connection, retrying with exponential backoff.
1002+
// primeWithRetry runs the configured ChannelPrimer with exponential backoff.
1003+
// Returns nil immediately when no primer is configured, so the pool can be
1004+
// used without priming.
10101005
func (cf *connectionFactory) primeWithRetry(ctx context.Context, conn *BigtableConn) error {
1006+
if cf.primer == nil {
1007+
return nil
1008+
}
10111009
backoffPolicy := gax.Backoff{
10121010
Initial: 100 * time.Millisecond,
10131011
Max: 2 * time.Second,
@@ -1022,7 +1020,7 @@ func (cf *connectionFactory) primeWithRetry(ctx context.Context, conn *BigtableC
10221020
return fmt.Errorf("bigtable_connpool: error before prime attempt %d: %w", attempt, err)
10231021
}
10241022

1025-
lastErr = conn.Prime(ctx, cf.instanceName, cf.appProfile, cf.featureFlagsMD)
1023+
lastErr = cf.primer.Prime(ctx, conn)
10261024
if lastErr == nil {
10271025
return nil
10281026
}

0 commit comments

Comments
 (0)