Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 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
6 changes: 3 additions & 3 deletions go/vt/vttablet/tabletserver/messager/message_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,17 +28,17 @@ import (
"golang.org/x/sync/semaphore"

"vitess.io/vitess/go/mysql/replication"

"vitess.io/vitess/go/sqltypes"
"vitess.io/vitess/go/stats"
"vitess.io/vitess/go/timer"
"vitess.io/vitess/go/vt/log"
binlogdatapb "vitess.io/vitess/go/vt/proto/binlogdata"
querypb "vitess.io/vitess/go/vt/proto/query"
"vitess.io/vitess/go/vt/sqlparser"
"vitess.io/vitess/go/vt/vttablet/tabletserver/schema"
"vitess.io/vitess/go/vt/vttablet/tabletserver/tabletenv"
"vitess.io/vitess/go/vt/vttablet/tabletserver/throttle/throttlerapp"

binlogdatapb "vitess.io/vitess/go/vt/proto/binlogdata"
querypb "vitess.io/vitess/go/vt/proto/query"
)

var (
Expand Down
64 changes: 56 additions & 8 deletions go/vt/vttablet/tabletserver/tabletserver.go
Original file line number Diff line number Diff line change
Expand Up @@ -551,7 +551,10 @@ func (tsv *TabletServer) begin(ctx context.Context, target *querypb.Target, save
logStats.OriginalSQL = beginSQL
if beginSQL != "" {
tsv.stats.QueryTimings.Record("BEGIN", startTime)
tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), startTime)
// With a tabletenv.LocalContext() the target can be nil.
if target != nil {
tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), startTime)
}
} else {
logStats.Method = ""
}
Expand Down Expand Up @@ -585,6 +588,23 @@ func (tsv *TabletServer) getPriorityFromOptions(options *querypb.ExecuteOptions)
return optionsPriority
}

// resolveTargetType returns the appropriate target tablet type for a
// TabletServer request. If the caller has a local context then it's
// an internal request and the target is the local tablet's current
// target.
func (tsv *TabletServer) resolveTargetType(ctx context.Context, target *querypb.Target) (string, error) {
if target != nil {
return target.TabletType.String(), nil
}
if !tabletenv.IsLocalContext(ctx) {
return topodatapb.ShardReplicationError_UNKNOWN.String(), vterrors.Errorf(vtrpcpb.Code_INVALID_ARGUMENT, "no target specified")
}
if tsv.sm.Target() == nil {
return topodatapb.ShardReplicationError_UNKNOWN.String(), vterrors.Errorf(vtrpcpb.Code_FAILED_PRECONDITION, "TabletServer has no current target")
}
return tsv.sm.Target().String(), nil

@deepthi deepthi Dec 8, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Interestingly:
https://github.com/vitessio/vitess/blob/main/go/vt/vttablet/tabletserver/tabletserver.go#L168-L172
Should we add a nil check on tsv.sm.Target() in line 172 while we are at it?

}

// Commit commits the specified transaction.
func (tsv *TabletServer) Commit(ctx context.Context, target *querypb.Target, transactionID int64) (newReservedID int64, err error) {
err = tsv.execRequest(
Expand All @@ -607,7 +627,11 @@ func (tsv *TabletServer) Commit(ctx context.Context, target *querypb.Target, tra
// handlePanicAndSendLogStats doesn't log the no-op.
if commitSQL != "" {
tsv.stats.QueryTimings.Record("COMMIT", startTime)
tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), startTime)
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, startTime)
} else {
logStats.Method = ""
}
Expand All @@ -625,7 +649,11 @@ func (tsv *TabletServer) Rollback(ctx context.Context, target *querypb.Target, t
target, nil, true, /* allowOnShutdown */
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("ROLLBACK", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
defer tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, time.Now())
logStats.TransactionID = transactionID
newReservedID, err = tsv.te.Rollback(ctx, transactionID)
if newReservedID > 0 {
Expand Down Expand Up @@ -1240,7 +1268,11 @@ func (tsv *TabletServer) ReserveBeginExecute(ctx context.Context, target *queryp
target, options, false, /* allowOnShutdown */
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("RESERVE", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
defer tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, time.Now())
connID, sessionStateChanges, err = tsv.te.ReserveBegin(ctx, options, preQueries, postBeginQueries)
if err != nil {
return err
Expand Down Expand Up @@ -1286,7 +1318,11 @@ func (tsv *TabletServer) ReserveBeginStreamExecute(
target, options, false, /* allowOnShutdown */
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("RESERVE", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
defer tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, time.Now())
connID, sessionStateChanges, err = tsv.te.ReserveBegin(ctx, options, preQueries, postBeginQueries)
if err != nil {
return err
Expand Down Expand Up @@ -1340,7 +1376,11 @@ func (tsv *TabletServer) ReserveExecute(ctx context.Context, target *querypb.Tar
target, options, allowOnShutdown,
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("RESERVE", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
defer tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, time.Now())
state.ReservedID, err = tsv.te.Reserve(ctx, options, transactionID, preQueries)
if err != nil {
return err
Expand Down Expand Up @@ -1391,7 +1431,11 @@ func (tsv *TabletServer) ReserveStreamExecute(
target, options, allowOnShutdown,
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("RESERVE", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
tabletTypeStr, err := tsv.resolveTargetType(ctx, target)
if err != nil {
return err
}
defer tsv.stats.QueryTimingsByTabletType.Record(tabletTypeStr, time.Now())
state.ReservedID, err = tsv.te.Reserve(ctx, options, transactionID, preQueries)
if err != nil {
return err
Expand Down Expand Up @@ -1421,7 +1465,10 @@ func (tsv *TabletServer) Release(ctx context.Context, target *querypb.Target, tr
target, nil, true, /* allowOnShutdown */
func(ctx context.Context, logStats *tabletenv.LogStats) error {
defer tsv.stats.QueryTimings.Record("RELEASE", time.Now())
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())

if target != nil {
defer tsv.stats.QueryTimingsByTabletType.Record(target.TabletType.String(), time.Now())
}
logStats.TransactionID = transactionID
logStats.ReservedID = reservedID
if reservedID != 0 {
Expand Down Expand Up @@ -1505,6 +1552,7 @@ func (tsv *TabletServer) execRequest(
span.Annotate("workload_name", options.WorkloadName)
}
trace.AnnotateSQL(span, sqlparser.Preview(sql))
// With a tabletenv.LocalContext() the target will be nil.
if target != nil {
span.Annotate("cell", target.Cell)
span.Annotate("shard", target.Shard)
Expand Down
27 changes: 27 additions & 0 deletions go/vt/vttablet/tabletserver/tabletserver_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -566,6 +566,33 @@ func TestTabletServerCommitPrepared(t *testing.T) {
require.NoError(t, err)
}

func TestTabletServerWithNilTarget(t *testing.T) {
// A non-nil target is required when not using a local context.
ctx := tabletenv.LocalContext()
db, tsv := setupTabletServerTest(t, ctx, "")
defer tsv.StopService()
defer db.Close()

executeSQL := "select * from test_table limit 1000"
executeSQLResult := &sqltypes.Result{
Fields: []*querypb.Field{
{Type: sqltypes.VarBinary},
},
Rows: [][]sqltypes.Value{
{sqltypes.NewVarBinary("row01")},
},
}
db.AddQuery(executeSQL, executeSQLResult)

target := (*querypb.Target)(nil)
state, err := tsv.Begin(ctx, target, nil)
require.NoError(t, err)
_, err = tsv.Execute(ctx, target, executeSQL, nil, state.TransactionID, 0, nil)
require.NoError(t, err)
_, err = tsv.Rollback(ctx, target, state.TransactionID)
require.NoError(t, err)
}

func TestSmallerTimeout(t *testing.T) {
testcases := []struct {
t1, t2, want time.Duration
Expand Down