Skip to content

Commit 95a1721

Browse files
bgwinesclaude
andauthored
Gate vttablet conn pools with Snake load shedder (#864)
* Gate OLTP read pool with Snake load shedder Wire the CoDel-based Snake gate into QueryExecutor.getConn() so that OLTP read requests are admitted or shed before consuming a MySQL connection. Controlled by --snake-enabled flag (off by default). ContentionID is extracted from the unique_id SQL margin comment that webapp injects, with a random 16-char hex fallback for requests without one. AI disclosure: Claude Code assisted with development. Every line of code was either written by or carefully reviewed by me :) Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> Signed-off-by: Brett Wines <bwines@slack-corp.com> * Use proto UniqueId as Snake contentionID, bypass valve when unset Replace the extractUniqueID regex (which targeted a comment format webapp never sends) with the proto-based approach: read qre.options.GetUniqueId() and only enter the Snake gate when it's non-empty. Requests without a unique_id bypass the gate entirely rather than sharing a single valve slot or using random IDs. Add `string unique_id = 22` to ExecuteOptions in proto/query.proto. This is an optional field — callers that don't set it get zero-cost bypass. The sqlparser/vtgate propagation to populate it from SQL directives is deferred to a follow-up PR. AI disclosure: Claude Code assisted with development. Every line of code was either written by or carefully reviewed by me :) Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> Signed-off-by: Brett Wines <bwines@slack-corp.com> * Gate DML with Snake load shedder (#871) Acquire a Snake slot in TxPool.Begin (for fresh connections only, not reserved conns) and release it in txComplete. This gates both autocommit and explicit transaction DML through the same CoDel-based load-shedding mechanism already used for the OLTP read pool. Requests without a unique_id bypass the gate entirely, matching the OLTP-read behavior. AI disclosure: Claude Code assisted with development. Every line of code was either written by or carefully reviewed by me :) Signed-off-by: Brett Wines <bwines@slack-corp.com> Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com> --------- Signed-off-by: Brett Wines <bwines@slack-corp.com> Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 2b970c6 commit 95a1721

11 files changed

Lines changed: 2680 additions & 1329 deletions

File tree

go/vt/proto/query/query.pb.go

Lines changed: 546 additions & 975 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

go/vt/proto/query/query_vtproto.pb.go

Lines changed: 1678 additions & 346 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

go/vt/vttablet/tabletserver/query_engine.go

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ import (
4646
tacl "vitess.io/vitess/go/vt/tableacl/acl"
4747
"vitess.io/vitess/go/vt/vterrors"
4848
"vitess.io/vitess/go/vt/vttablet/tabletserver/connpool"
49+
"vitess.io/vitess/go/vt/vttablet/tabletserver/loadshed"
4950
"vitess.io/vitess/go/vt/vttablet/tabletserver/planbuilder"
5051
"vitess.io/vitess/go/vt/vttablet/tabletserver/rules"
5152
"vitess.io/vitess/go/vt/vttablet/tabletserver/schema"
@@ -197,6 +198,9 @@ type QueryEngine struct {
197198
accessCheckerLogger *logutil.ThrottledLogger
198199

199200
redactUIQuery bool
201+
202+
// snake is the CoDel-based load-shedding gate for the OLTP read pool.
203+
snake *loadshed.Snake
200204
}
201205

202206
// NewQueryEngine creates a new QueryEngine.
@@ -241,6 +245,20 @@ func NewQueryEngine(env tabletenv.Env, se *schema.Engine) *QueryEngine {
241245
}
242246
qe.txSerializer = txserializer.New(env)
243247

248+
if config.SnakeEnabled {
249+
qe.snake = loadshed.NewSnake(loadshed.SnakeConfig{
250+
Name: "oltp-read",
251+
CoDel: loadshed.CoDelConfig{
252+
TargetNs: func() int64 { return config.SnakeTarget.Nanoseconds() },
253+
IntervalNs: func() int64 { return config.SnakeInterval.Nanoseconds() },
254+
Exponent: func() float64 { return 0.5 },
255+
MinDropDelayNs: func() int64 { return int64(time.Millisecond) },
256+
},
257+
Capacity: func() int { return config.OltpReadPool.Size },
258+
LoadsheddingAllowed: func() bool { return true },
259+
})
260+
}
261+
244262
qe.strictTableACL = config.StrictTableACL
245263
qe.enableTableACLDryRun = config.EnableTableACLDryRun
246264

go/vt/vttablet/tabletserver/query_executor.go

Lines changed: 32 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -714,12 +714,13 @@ func (qre *QueryExecutor) execSelect() (*sqltypes.Result, error) {
714714

715715
if original {
716716
defer q.Broadcast()
717-
conn, err := qre.getConn()
717+
conn, release, err := qre.getConn()
718718

719719
if err != nil {
720720
q.SetErr(err)
721721
} else {
722722
defer conn.Recycle()
723+
defer release()
723724
res, err := qre.execDBConn(conn.Conn, sql, true)
724725
if qre.tsv.config.ConsolidatorCacheProto3Rows && q.HasWaiters() {
725726
res.CacheProto3Rows()
@@ -755,11 +756,12 @@ func (qre *QueryExecutor) execSelect() (*sqltypes.Result, error) {
755756
}
756757
// If waiter cap exceeded, fall through to independent execution
757758
}
758-
conn, err := qre.getConn()
759+
conn, release, err := qre.getConn()
759760
if err != nil {
760761
return nil, err
761762
}
762763
defer conn.Recycle()
764+
defer release()
763765
res, err := qre.execDBConn(conn.Conn, sql, true)
764766
if err != nil {
765767
return nil, err
@@ -797,22 +799,45 @@ func (qre *QueryExecutor) verifyRowCount(count, maxrows int64) error {
797799
}
798800

799801
func (qre *QueryExecutor) execOther() (*sqltypes.Result, error) {
800-
conn, err := qre.getConn()
802+
conn, release, err := qre.getConn()
801803
if err != nil {
802804
return nil, err
803805
}
804806
defer conn.Recycle()
807+
defer release()
805808
return qre.execDBConn(conn.Conn, qre.query, true)
806809
}
807810

808-
func (qre *QueryExecutor) getConn() (*connpool.PooledConn, error) {
811+
func (qre *QueryExecutor) getConn() (*connpool.PooledConn, func(), error) {
809812
span, ctx := trace.NewSpan(qre.ctx, "QueryExecutor.getConn")
810813
defer span.Finish()
811814

812815
defer func(start time.Time) {
813816
qre.logStats.WaitingForConnection += time.Since(start)
814817
}(time.Now())
815-
return qre.tsv.qe.conns.Get(ctx, qre.setting)
818+
819+
snake := qre.tsv.qe.snake
820+
if snake != nil {
821+
contentionID := qre.options.GetUniqueId()
822+
if contentionID != "" {
823+
unlock, err := snake.Acquire(ctx, contentionID)
824+
if err != nil {
825+
return nil, nil, vterrors.Errorf(vtrpcpb.Code_RESOURCE_EXHAUSTED, "load shed: %v", err)
826+
}
827+
conn, err := qre.tsv.qe.conns.Get(ctx, qre.setting)
828+
if err != nil {
829+
unlock.Release()
830+
return nil, nil, err
831+
}
832+
return conn, func() { unlock.Release() }, nil
833+
}
834+
}
835+
836+
conn, err := qre.tsv.qe.conns.Get(ctx, qre.setting)
837+
if err != nil {
838+
return nil, nil, err
839+
}
840+
return conn, func() {}, nil
816841
}
817842

818843
func (qre *QueryExecutor) getStreamConn() (*connpool.PooledConn, error) {
@@ -927,11 +952,12 @@ func rewriteOUTParamError(err error) error {
927952
}
928953

929954
func (qre *QueryExecutor) execCallProc() (*sqltypes.Result, error) {
930-
conn, err := qre.getConn()
955+
conn, release, err := qre.getConn()
931956
if err != nil {
932957
return nil, err
933958
}
934959
defer conn.Recycle()
960+
defer release()
935961
sql, _, err := qre.generateFinalSQL(qre.plan.FullQuery, qre.bindVars)
936962
if err != nil {
937963
return nil, err

go/vt/vttablet/tabletserver/query_executor_test.go

Lines changed: 89 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1621,10 +1621,11 @@ func TestGetConnectionLogStats(t *testing.T) {
16211621

16221622
// getConn() happy path
16231623
qre := newTestQueryExecutor(ctx, tsv, input, 0)
1624-
conn, err := qre.getConn()
1624+
conn, release, err := qre.getConn()
16251625
assert.NoError(t, err)
16261626
assert.NotNil(t, conn)
16271627
assert.True(t, qre.logStats.WaitingForConnection > 0)
1628+
release()
16281629

16291630
// getStreamConn() happy path
16301631
qre = newTestQueryExecutor(ctx, tsv, input, 0)
@@ -1638,7 +1639,7 @@ func TestGetConnectionLogStats(t *testing.T) {
16381639

16391640
// getConn() error path
16401641
qre = newTestQueryExecutor(ctx, tsv, input, 0)
1641-
_, err = qre.getConn()
1642+
_, _, err = qre.getConn()
16421643
assert.Error(t, err)
16431644
assert.True(t, qre.logStats.WaitingForConnection > 0)
16441645

@@ -1649,6 +1650,92 @@ func TestGetConnectionLogStats(t *testing.T) {
16491650
assert.True(t, qre.logStats.WaitingForConnection > 0)
16501651
}
16511652

1653+
func TestGetConnSnakeBypassed(t *testing.T) {
1654+
db := setUpQueryExecutorTest(t)
1655+
defer db.Close()
1656+
1657+
ctx := context.Background()
1658+
cfg := tabletenv.NewDefaultConfig()
1659+
cfg.OltpReadPool.Size = 2
1660+
cfg.TxPool.Size = 100
1661+
cfg.SnakeEnabled = true
1662+
cfg.SnakeTarget = 5 * time.Millisecond
1663+
cfg.SnakeInterval = 100 * time.Millisecond
1664+
cfg.DB = newDBConfigs(db)
1665+
1666+
srvTopoCounts := stats.NewCountersWithSingleLabel("", "Resilient srvtopo server operations", "type")
1667+
tsv := NewTabletServer(ctx, vtenv.NewTestEnv(), "TabletServerTest", cfg, memorytopo.NewServer(ctx, ""), &topodatapb.TabletAlias{}, srvTopoCounts)
1668+
target := &querypb.Target{TabletType: topodatapb.TabletType_PRIMARY}
1669+
err := tsv.StartService(target, cfg.DB, nil)
1670+
require.NoError(t, err)
1671+
defer tsv.StopService()
1672+
1673+
require.NotNil(t, tsv.qe.snake)
1674+
1675+
input := "select * from test_table limit 1"
1676+
1677+
// Without UniqueId set, Snake gate is bypassed entirely
1678+
qre := newTestQueryExecutor(ctx, tsv, input, 0)
1679+
conn, release, err := qre.getConn()
1680+
require.NoError(t, err)
1681+
require.NotNil(t, conn)
1682+
conn.Recycle()
1683+
release()
1684+
}
1685+
1686+
func TestGetConnWithSnake(t *testing.T) {
1687+
db := setUpQueryExecutorTest(t)
1688+
defer db.Close()
1689+
1690+
ctx := context.Background()
1691+
cfg := tabletenv.NewDefaultConfig()
1692+
cfg.OltpReadPool.Size = 2
1693+
cfg.TxPool.Size = 100
1694+
cfg.SnakeEnabled = true
1695+
cfg.SnakeTarget = 5 * time.Millisecond
1696+
cfg.SnakeInterval = 100 * time.Millisecond
1697+
cfg.DB = newDBConfigs(db)
1698+
1699+
srvTopoCounts := stats.NewCountersWithSingleLabel("", "Resilient srvtopo server operations", "type")
1700+
tsv := NewTabletServer(ctx, vtenv.NewTestEnv(), "TabletServerTest", cfg, memorytopo.NewServer(ctx, ""), &topodatapb.TabletAlias{}, srvTopoCounts)
1701+
target := &querypb.Target{TabletType: topodatapb.TabletType_PRIMARY}
1702+
err := tsv.StartService(target, cfg.DB, nil)
1703+
require.NoError(t, err)
1704+
defer tsv.StopService()
1705+
1706+
require.NotNil(t, tsv.qe.snake, "snake should be initialized when SnakeEnabled=true")
1707+
1708+
input := "select * from test_table limit 1"
1709+
1710+
// With UniqueId set, Snake gate admits the request
1711+
qre := newTestQueryExecutor(ctx, tsv, input, 0)
1712+
qre.options = &querypb.ExecuteOptions{UniqueId: "test-request-123"}
1713+
conn, release, err := qre.getConn()
1714+
require.NoError(t, err)
1715+
require.NotNil(t, conn)
1716+
conn.Recycle()
1717+
release()
1718+
}
1719+
1720+
func TestGetConnSnakeDisabled(t *testing.T) {
1721+
db := setUpQueryExecutorTest(t)
1722+
defer db.Close()
1723+
1724+
ctx := context.Background()
1725+
tsv := newTestTabletServer(ctx, noFlags, db)
1726+
defer tsv.StopService()
1727+
1728+
assert.Nil(t, tsv.qe.snake, "snake should be nil when SnakeEnabled=false")
1729+
1730+
input := "select * from test_table limit 1"
1731+
qre := newTestQueryExecutor(ctx, tsv, input, 0)
1732+
conn, release, err := qre.getConn()
1733+
require.NoError(t, err)
1734+
require.NotNil(t, conn)
1735+
conn.Recycle()
1736+
release()
1737+
}
1738+
16521739
type executorFlags int64
16531740

16541741
const (

go/vt/vttablet/tabletserver/tabletenv/config.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,7 @@ var (
8080
unhealthyThreshold time.Duration
8181
transitionGracePeriod time.Duration
8282
enableReplicationReporter bool
83+
enableSnake bool
8384
)
8485

8586
func init() {
@@ -223,6 +224,10 @@ func registerTabletEnvFlags(fs *pflag.FlagSet) {
223224
fs.BoolVar(&currentConfig.EnablePerWorkloadTableMetrics, "enable-per-workload-table-metrics", defaultConfig.EnablePerWorkloadTableMetrics, "If true, query counts and query error metrics include a label that identifies the workload")
224225

225226
fs.BoolVar(&currentConfig.Unmanaged, "unmanaged", false, "Indicates an unmanaged tablet, i.e. using an external mysql-compatible database")
227+
228+
fs.BoolVar(&enableSnake, "snake-enabled", false, "If true, enables CoDel-based load shedding (Snake) on the OLTP read pool.")
229+
fs.DurationVar(&currentConfig.SnakeTarget, "snake-target", 5*time.Millisecond, "CoDel target delay for the Snake load shedder.")
230+
fs.DurationVar(&currentConfig.SnakeInterval, "snake-interval", 100*time.Millisecond, "CoDel interval for the Snake load shedder.")
226231
}
227232

228233
var (
@@ -295,6 +300,7 @@ func Init() {
295300
currentConfig.GracePeriods.Transition = transitionGracePeriod
296301
currentConfig.SemiSyncMonitor.Interval = semiSyncMonitorInterval
297302

303+
currentConfig.SnakeEnabled = enableSnake
298304
logFormat := streamlog.GetQueryLogConfig().Format
299305
switch logFormat {
300306
case streamlog.QueryLogFormatText:
@@ -387,6 +393,10 @@ type TabletConfig struct {
387393
EnableViews bool `json:"-"`
388394

389395
EnablePerWorkloadTableMetrics bool `json:"-"`
396+
397+
SnakeEnabled bool `json:"-"`
398+
SnakeTarget time.Duration `json:"-"`
399+
SnakeInterval time.Duration `json:"-"`
390400
}
391401

392402
func (cfg *TabletConfig) MarshalJSON() ([]byte, error) {

go/vt/vttablet/tabletserver/tx/api.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,8 @@ type (
5757
LogToFile bool
5858

5959
Stats *servenv.TimingsWrapper
60+
61+
SnakeRelease func()
6062
}
6163

6264
// Query contains the query and involved tables executed inside transaction.

go/vt/vttablet/tabletserver/tx_engine.go

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ import (
3434
"vitess.io/vitess/go/vt/sqlparser"
3535
"vitess.io/vitess/go/vt/vterrors"
3636
"vitess.io/vitess/go/vt/vttablet/tabletserver/connpool"
37+
"vitess.io/vitess/go/vt/vttablet/tabletserver/loadshed"
3738
"vitess.io/vitess/go/vt/vttablet/tabletserver/tabletenv"
3839
"vitess.io/vitess/go/vt/vttablet/tabletserver/tx"
3940
"vitess.io/vitess/go/vt/vttablet/tabletserver/txlimiter"
@@ -114,6 +115,19 @@ func NewTxEngine(env tabletenv.Env, dxNotifier func()) *TxEngine {
114115
}
115116
limiter := txlimiter.New(env)
116117
te.txPool = NewTxPool(env, limiter)
118+
if config.SnakeEnabled {
119+
te.txPool.snake = loadshed.NewSnake(loadshed.SnakeConfig{
120+
Name: "dml",
121+
CoDel: loadshed.CoDelConfig{
122+
TargetNs: func() int64 { return config.SnakeTarget.Nanoseconds() },
123+
IntervalNs: func() int64 { return config.SnakeInterval.Nanoseconds() },
124+
Exponent: func() float64 { return 0.5 },
125+
MinDropDelayNs: func() int64 { return int64(time.Millisecond) },
126+
},
127+
Capacity: func() int { return config.TxPool.Size },
128+
LoadsheddingAllowed: func() bool { return true },
129+
})
130+
}
117131
// We initially allow twoPC (handles vttablet restarts).
118132
// We will disallow them for a few reasons -
119133
// 1. When a new tablet is promoted if semi-sync is turned off.

go/vt/vttablet/tabletserver/tx_pool.go

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ import (
3131
"vitess.io/vitess/go/vt/servenv"
3232
"vitess.io/vitess/go/vt/sqlparser"
3333
"vitess.io/vitess/go/vt/vterrors"
34+
"vitess.io/vitess/go/vt/vttablet/tabletserver/loadshed"
3435
"vitess.io/vitess/go/vt/vttablet/tabletserver/tabletenv"
3536
"vitess.io/vitess/go/vt/vttablet/tabletserver/tx"
3637
"vitess.io/vitess/go/vt/vttablet/tabletserver/txlimiter"
@@ -68,6 +69,7 @@ type (
6869
scp *StatefulConnectionPool
6970
ticks *timer.Timer
7071
limiter txlimiter.TxLimiter
72+
snake *loadshed.Snake
7173

7274
logMu sync.Mutex
7375
lastLog time.Time
@@ -235,6 +237,7 @@ func (tp *TxPool) Begin(ctx context.Context, options *querypb.ExecuteOptions, re
235237

236238
var conn *StatefulConnection
237239
var err error
240+
var snakeRelease func()
238241
if reservedID != 0 {
239242
conn, err = tp.scp.GetAndLock(reservedID, "start transaction on reserve conn")
240243
if err != nil {
@@ -249,12 +252,27 @@ func (tp *TxPool) Begin(ctx context.Context, options *querypb.ExecuteOptions, re
249252
if !tp.limiter.Get(immediateCaller, effectiveCaller) {
250253
return nil, "", "", vterrors.Errorf(vtrpcpb.Code_RESOURCE_EXHAUSTED, "per-user transaction pool connection limit exceeded")
251254
}
255+
256+
if tp.snake != nil {
257+
if contentionID := options.GetUniqueId(); contentionID != "" {
258+
unlock, snakeErr := tp.snake.Acquire(ctx, contentionID)
259+
if snakeErr != nil {
260+
tp.limiter.Release(immediateCaller, effectiveCaller)
261+
return nil, "", "", vterrors.Errorf(vtrpcpb.Code_RESOURCE_EXHAUSTED, "dml load shed: %v", snakeErr)
262+
}
263+
snakeRelease = func() { unlock.Release() }
264+
}
265+
}
266+
252267
conn, err = tp.createConn(ctx, options, setting)
253268
defer func() {
254269
if err != nil {
255270
// The transaction limiter frees transactions on rollback or commit. If we fail to create the transaction,
256271
// release immediately since there will be no rollback or commit.
257272
tp.limiter.Release(immediateCaller, effectiveCaller)
273+
if snakeRelease != nil {
274+
snakeRelease()
275+
}
258276
}
259277
}()
260278
}
@@ -272,6 +290,7 @@ func (tp *TxPool) Begin(ctx context.Context, options *querypb.ExecuteOptions, re
272290
if setting != nil {
273291
conn.TxProperties().RecordQueryDetail(setting.ApplyQuery(), nil)
274292
}
293+
conn.TxProperties().SnakeRelease = snakeRelease
275294
return conn, sql, sessionStateChanges, nil
276295
}
277296

@@ -444,6 +463,9 @@ func (tp *TxPool) LogActive() {
444463

445464
func (tp *TxPool) txComplete(conn *StatefulConnection, reason tx.ReleaseReason) {
446465
conn.LogTransaction(reason)
466+
if release := conn.TxProperties().SnakeRelease; release != nil {
467+
release()
468+
}
447469
tp.limiter.Release(conn.TxProperties().ImmediateCaller, conn.TxProperties().EffectiveCaller)
448470
conn.CleanTxState()
449471
}

0 commit comments

Comments
 (0)