From 7773282d67f21e7a06da436c12b2e35f9d03b312 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Sun, 5 Jul 2026 22:06:38 +0800 Subject: [PATCH 01/16] starter ttl worker Signed-off-by: ystaticy --- cmd/tidb-server/main.go | 25 +++--- pkg/ddl/BUILD.bazel | 1 + pkg/ddl/create_table.go | 4 + pkg/ddl/ddl.go | 3 + pkg/ddl/options.go | 9 +++ pkg/ddl/ttl.go | 43 ++++++++++ pkg/ddl/ttl_test.go | 108 ++++++++++++++++++++++++++ pkg/domain/BUILD.bazel | 4 +- pkg/domain/domain.go | 28 ++++++- pkg/domain/domain_sysvars.go | 9 +++ pkg/domain/domain_test.go | 70 +++++++++++++++++ pkg/session/BUILD.bazel | 1 + pkg/session/bootstrap_test.go | 2 +- pkg/session/session.go | 18 ++++- pkg/session/session_test.go | 2 +- pkg/session/tidb.go | 12 ++- pkg/sessionctx/variable/sysvar.go | 6 +- pkg/sessionctx/variable/tidb_vars.go | 2 + pkg/ttl/ttlworker/BUILD.bazel | 4 + pkg/ttl/ttlworker/job_manager.go | 47 ++++++++++- pkg/ttl/ttlworker/job_manager_test.go | 96 +++++++++++++++++++++++ 21 files changed, 471 insertions(+), 23 deletions(-) diff --git a/cmd/tidb-server/main.go b/cmd/tidb-server/main.go index e117e3e138629..47f0cd2d511b1 100644 --- a/cmd/tidb-server/main.go +++ b/cmd/tidb-server/main.go @@ -477,11 +477,10 @@ func main() { keyspaceName := keyspace.GetKeyspaceNameBySettings() executor.Start() resourcemanager.InstanceResourceManager.Start() - storage, dom, err := createStoreDDLOwnerMgrAndDomain(keyspaceName) + storage, dom, externalWorkloadManager, err := createStoreDDLOwnerMgrAndDomain(keyspaceName) terror.MustNil(err) repository.SetupRepository(dom) - if deploymode.IsStarter() { - externalWorkloadManager := initExternalWorkloadManager(context.Background(), storage) + if externalWorkloadManager != nil { defer closeExternalWorkloadManager(externalWorkloadManager) } svr := createServer(storage, dom) @@ -611,7 +610,7 @@ func registerStores() error { return err } -func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.Domain, error) { +func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.Domain, extworkload.Manager, error) { if config.GetGlobalConfig().Store == config.StoreTypeUniStore { kv.StandAloneTiDB = true } @@ -622,27 +621,33 @@ func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.D if pdhttpCli != nil { pdStatus, err := pdhttpCli.GetStatus(context.Background()) if err != nil { - return nil, nil, err + return nil, nil, nil, err } if !kerneltype.IsMatch(pdStatus.KernelType) { log.Error("kernel type mismatch", zap.String("pd", pdStatus.KernelType), zap.String("tidb", kerneltype.Name())) - return nil, nil, errors.New("kernel type mismatch") + return nil, nil, nil, errors.New("kernel type mismatch") } } } + var externalWorkloadManager extworkload.Manager + if deploymode.IsStarter() { + externalWorkloadManager = initExternalWorkloadManager(context.Background(), storage) + } copr.GlobalMPPFailedStoreProber.Run() mppcoordmanager.InstanceMPPCoordinatorManager.Run() // Bootstrap a session to load information schema. err := ddl.StartOwnerManager(context.Background(), storage) if err != nil { - return nil, nil, err + closeExternalWorkloadManager(externalWorkloadManager) + return nil, nil, nil, err } - dom, err := session.BootstrapSession(storage) + dom, err := session.BootstrapSessionWithExternalWorkloadManager(storage, externalWorkloadManager) if err != nil { - return nil, nil, err + closeExternalWorkloadManager(externalWorkloadManager) + return nil, nil, nil, err } - return storage, dom, nil + return storage, dom, externalWorkloadManager, nil } // Prometheus push. diff --git a/pkg/ddl/BUILD.bazel b/pkg/ddl/BUILD.bazel index 801669888844e..547c5da2fbadd 100644 --- a/pkg/ddl/BUILD.bazel +++ b/pkg/ddl/BUILD.bazel @@ -112,6 +112,7 @@ go_library( "//pkg/expression", "//pkg/expression/exprctx", "//pkg/expression/exprstatic", + "//pkg/extworkload", "//pkg/infoschema", "//pkg/infoschema/context", "//pkg/ingestor/engineapi", diff --git a/pkg/ddl/create_table.go b/pkg/ddl/create_table.go index b76b020872c90..689dadd75856a 100644 --- a/pkg/ddl/create_table.go +++ b/pkg/ddl/create_table.go @@ -249,6 +249,10 @@ func (w *worker) onCreateTable(jobCtx *jobContext, job *model.Job) (ver int64, _ return ver, errors.Trace(err) } + if err := w.registerTTLTableToExternalWorkload(jobCtx.ctx, tbInfo); err != nil { + return ver, errors.Trace(err) + } + // Finish this job. job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo) return ver, errors.Trace(err) diff --git a/pkg/ddl/ddl.go b/pkg/ddl/ddl.go index 09e65ea12455d..f1111a5a714a3 100644 --- a/pkg/ddl/ddl.go +++ b/pkg/ddl/ddl.go @@ -46,6 +46,7 @@ import ( "github.com/pingcap/tidb/pkg/dxf/framework/proto" "github.com/pingcap/tidb/pkg/dxf/framework/scheduler" "github.com/pingcap/tidb/pkg/dxf/framework/taskexecutor" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" "github.com/pingcap/tidb/pkg/kv" "github.com/pingcap/tidb/pkg/meta" @@ -348,6 +349,7 @@ type ddlCtx struct { etcdCli *clientv3.Client autoidCli *autoid.ClientDiscover schemaLoader SchemaLoader + extWorkload extworkload.Manager // reorgCtx is used for reorganization. reorgCtx reorgContexts @@ -780,6 +782,7 @@ func newDDL(ctx context.Context, options ...Option) (*ddl, *executor) { etcdCli: opt.EtcdCli, autoidCli: opt.AutoIDClient, schemaLoader: opt.SchemaLoader, + extWorkload: opt.ExtWorkloadMgr, } ddlCtx.reorgCtx.reorgCtxMap = make(map[int64]*reorgCtx) ddlCtx.jobCtx.jobCtxMap = make(map[int64]*ReorgContext) diff --git a/pkg/ddl/options.go b/pkg/ddl/options.go index 7bd683551c5ae..9fbc3e191d350 100644 --- a/pkg/ddl/options.go +++ b/pkg/ddl/options.go @@ -18,6 +18,7 @@ import ( "time" "github.com/pingcap/tidb/pkg/ddl/notifier" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" "github.com/pingcap/tidb/pkg/kv" "github.com/pingcap/tidb/pkg/meta/autoid" @@ -36,6 +37,7 @@ type Options struct { Lease time.Duration SchemaLoader SchemaLoader EventPublishStore notifier.Store + ExtWorkloadMgr extworkload.Manager } // WithEtcdClient specifies the `clientv3.Client` of DDL used to request the etcd service @@ -86,3 +88,10 @@ func WithEventPublishStore(store notifier.Store) Option { options.EventPublishStore = store } } + +// WithExternalWorkloadManager specifies the manager used to coordinate external background workloads. +func WithExternalWorkloadManager(manager extworkload.Manager) Option { + return func(options *Options) { + options.ExtWorkloadMgr = manager + } +} diff --git a/pkg/ddl/ttl.go b/pkg/ddl/ttl.go index c8ac853c73675..f3adfee3c51ca 100644 --- a/pkg/ddl/ttl.go +++ b/pkg/ddl/ttl.go @@ -15,18 +15,23 @@ package ddl import ( + "context" "strings" "time" "github.com/pingcap/errors" + "github.com/pingcap/tidb/pkg/ddl/logutil" + "github.com/pingcap/tidb/pkg/extworkload" infoschemactx "github.com/pingcap/tidb/pkg/infoschema/context" "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/parser/ast" "github.com/pingcap/tidb/pkg/parser/format" "github.com/pingcap/tidb/pkg/parser/mysql" + "github.com/pingcap/tidb/pkg/sessionctx/vardef" "github.com/pingcap/tidb/pkg/ttl/cache" "github.com/pingcap/tidb/pkg/types" "github.com/pingcap/tidb/pkg/util/dbterror" + "go.uber.org/zap" ) func onTTLInfoRemove(jobCtx *jobContext, job *model.Job) (ver int64, err error) { @@ -40,6 +45,9 @@ func onTTLInfoRemove(jobCtx *jobContext, job *model.Job) (ver int64, err error) if err != nil { return ver, errors.Trace(err) } + if jobCtx.oldDDLCtx != nil { + jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID) + } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) return ver, nil } @@ -87,10 +95,45 @@ func onTTLInfoChange(jobCtx *jobContext, job *model.Job) (ver int64, err error) if err != nil { return ver, errors.Trace(err) } + if jobCtx.oldDDLCtx != nil { + if err := jobCtx.oldDDLCtx.registerTTLTableToExternalWorkload(jobCtx.ctx, tblInfo); err != nil { + logutil.DDLLogger().Warn("failed to register TTL table to external workload controller", + zap.Int64("tableID", tblInfo.ID), + zap.Error(err)) + } + } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) return ver, nil } +func (dc *ddlCtx) externalWorkloadMaster() (extworkload.Manager, bool) { + if dc == nil { + return nil, false + } + manager := dc.extWorkload + return manager, extworkload.IsMaster(manager) +} + +func (dc *ddlCtx) registerTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) error { + manager, ok := dc.externalWorkloadMaster() + if !ok || tblInfo == nil || tblInfo.TTLInfo == nil || !tblInfo.TTLInfo.Enable { + return nil + } + return manager.RegisterTTLTask(ctx, tblInfo.ID, vardef.EnableTTLJob.Load()) +} + +func (dc *ddlCtx) deleteTTLTableFromExternalWorkload(ctx context.Context, tableID int64) { + manager, ok := dc.externalWorkloadMaster() + if !ok { + return + } + if err := manager.DeleteTTLTableInfo(ctx, tableID); err != nil { + logutil.DDLLogger().Warn("failed to delete TTL table from external workload controller", + zap.Int64("tableID", tableID), + zap.Error(err)) + } +} + // checkTTLInfoValid checks the TTL settings for a table. // The argument `isForForeignKeyCheck` is used to check the table should not be referenced by foreign key. // If `isForForeignKeyCheck` is `nil`, it will skip the foreign key check. diff --git a/pkg/ddl/ttl_test.go b/pkg/ddl/ttl_test.go index 323a306749914..72e05ee1ba7dd 100644 --- a/pkg/ddl/ttl_test.go +++ b/pkg/ddl/ttl_test.go @@ -15,11 +15,17 @@ package ddl import ( + "context" + "errors" "testing" + "github.com/pingcap/kvproto/pkg/keyspacepb" + "github.com/pingcap/tidb/pkg/config" "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/parser/ast" + "github.com/pingcap/tidb/pkg/sessionctx/vardef" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) func Test_getTTLInfoInOptions(t *testing.T) { @@ -148,3 +154,105 @@ func Test_getTTLInfoInOptions(t *testing.T) { assert.Equal(t, c.err, err) } } + +type fakeExternalWorkloadManager struct { + role config.ExternalWorkloadRole + registeredTable int64 + registerEnabled bool + deletedTable int64 + recycledCreateTS uint64 + updatedEnable *bool + registerErr error +} + +func (m *fakeExternalWorkloadManager) Close() error { return nil } +func (m *fakeExternalWorkloadManager) Role() config.ExternalWorkloadRole { + return m.role +} +func (*fakeExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil } +func (*fakeExternalWorkloadManager) InitializeGCV2(context.Context) error { + return nil +} +func (*fakeExternalWorkloadManager) AbortGCV2(context.Context) error { return nil } +func (*fakeExternalWorkloadManager) RegisterGCV2(context.Context, uint64, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) UpdateGCLifeTime(context.Context, int64) error { + return nil +} +func (m *fakeExternalWorkloadManager) RegisterTTLTask(_ context.Context, tableID int64, ttlJobEnable bool) error { + m.registeredTable = tableID + m.registerEnabled = ttlJobEnable + return m.registerErr +} +func (m *fakeExternalWorkloadManager) DeleteTTLTableInfo(_ context.Context, tableID int64) error { + m.deletedTable = tableID + return nil +} +func (m *fakeExternalWorkloadManager) RecycleTTLTask(_ context.Context, completedJobCreateTime uint64) error { + m.recycledCreateTS = completedJobCreateTime + return nil +} +func (m *fakeExternalWorkloadManager) UpdateTTLJobEnable(_ context.Context, ttlJobEnable bool) error { + m.updatedEnable = &ttlJobEnable + return nil +} +func (*fakeExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error { + return nil +} + +func TestExternalWorkloadTTLTableReportsOnlyFromMaster(t *testing.T) { + origEnable := vardef.EnableTTLJob.Load() + vardef.EnableTTLJob.Store(false) + defer vardef.EnableTTLJob.Store(origEnable) + + tblInfo := &model.TableInfo{ + ID: 123, + TTLInfo: &model.TTLInfo{ + Enable: true, + }, + } + + master := &fakeExternalWorkloadManager{role: config.RoleMaster} + dc := &ddlCtx{extWorkload: master} + require.NoError(t, dc.registerTTLTableToExternalWorkload(context.Background(), tblInfo)) + require.Equal(t, int64(123), master.registeredTable) + require.False(t, master.registerEnabled) + + dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID) + require.Equal(t, int64(123), master.deletedTable) + + ttlWorker := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} + dc = &ddlCtx{extWorkload: ttlWorker} + require.NoError(t, dc.registerTTLTableToExternalWorkload(context.Background(), tblInfo)) + dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID) + require.Zero(t, ttlWorker.registeredTable) + require.Zero(t, ttlWorker.deletedTable) +} + +func TestExternalWorkloadTTLTableRegisterSkipsDisabledTTL(t *testing.T) { + manager := &fakeExternalWorkloadManager{role: config.RoleMaster} + dc := &ddlCtx{extWorkload: manager} + require.NoError(t, dc.registerTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ + ID: 123, + TTLInfo: &model.TTLInfo{Enable: false}, + })) + require.Zero(t, manager.registeredTable) +} + +func TestExternalWorkloadTTLTableRegisterReturnsError(t *testing.T) { + boom := errors.New("boom") + manager := &fakeExternalWorkloadManager{role: config.RoleMaster, registerErr: boom} + dc := &ddlCtx{extWorkload: manager} + err := dc.registerTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ + ID: 123, + TTLInfo: &model.TTLInfo{Enable: true}, + }) + require.ErrorIs(t, err, boom) +} diff --git a/pkg/domain/BUILD.bazel b/pkg/domain/BUILD.bazel index 3b18f8d2dce0f..7644c96fae844 100644 --- a/pkg/domain/BUILD.bazel +++ b/pkg/domain/BUILD.bazel @@ -43,6 +43,7 @@ go_library( "//pkg/dxf/framework/storage", "//pkg/dxf/framework/taskexecutor", "//pkg/errno", + "//pkg/extworkload", "//pkg/infoschema", "//pkg/infoschema/issyncer", "//pkg/infoschema/isvalidator", @@ -154,7 +155,7 @@ go_test( ], embed = [":domain"], flaky = True, - shard_count = 30, + shard_count = 32, deps = [ "//pkg/config", "//pkg/ddl", @@ -191,6 +192,7 @@ go_test( "@com_github_ngaut_pools//:pools", "@com_github_pingcap_errors//:errors", "@com_github_pingcap_failpoint//:failpoint", + "@com_github_pingcap_kvproto//pkg/keyspacepb", "@com_github_pingcap_kvproto//pkg/metapb", "@com_github_pingcap_kvproto//pkg/resource_manager", "@com_github_prometheus_client_model//go", diff --git a/pkg/domain/domain.go b/pkg/domain/domain.go index 381b8f6a4bf6e..66cfca1fe533f 100644 --- a/pkg/domain/domain.go +++ b/pkg/domain/domain.go @@ -56,6 +56,7 @@ import ( "github.com/pingcap/tidb/pkg/dxf/framework/storage" "github.com/pingcap/tidb/pkg/dxf/framework/taskexecutor" "github.com/pingcap/tidb/pkg/errno" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" "github.com/pingcap/tidb/pkg/infoschema/issyncer" "github.com/pingcap/tidb/pkg/infoschema/isvalidator" @@ -234,6 +235,7 @@ type Domain struct { // only used for nextgen crossKSSessMgr *crossks.Manager crossKSSessFactoryGetter func(string, validatorapi.Validator) pools.Factory + extWorkloadMgr extworkload.Manager } var _ sqlsvrapi.Server = (*Domain)(nil) @@ -700,6 +702,7 @@ func (do *Domain) Init( ddl.WithLease(do.schemaLease), ddl.WithSchemaLoader(do.isSyncer), ddl.WithEventPublishStore(ddlNotifierStore), + ddl.WithExternalWorkloadManager(do.extWorkloadMgr), ) failpoint.Inject("MockReplaceDDL", func(val failpoint.Value) { @@ -937,6 +940,17 @@ func (do *Domain) SetOnClose(onClose func()) { do.onClose = onClose } +// SetExternalWorkloadManager installs the external workload manager used by +// background components that coordinate work with an external controller. +func (do *Domain) SetExternalWorkloadManager(manager extworkload.Manager) { + do.extWorkloadMgr = manager +} + +// ExternalWorkloadManager returns the external workload manager for this domain. +func (do *Domain) ExternalWorkloadManager() extworkload.Manager { + return do.extWorkloadMgr +} + func (do *Domain) initLogBackup(ctx context.Context, pdClient pd.Client) error { cfg := config.GetGlobalConfig() if pdClient == nil || do.etcdClient == nil { @@ -2875,11 +2889,23 @@ func (do *Domain) serverIDKeeper() { // StartTTLJobManager creates and starts the ttl job manager func (do *Domain) StartTTLJobManager() { - ttlJobManager := ttlworker.NewJobManager(do.ddl.GetID(), do.advancedSysSessionPool, do.store, do.etcdClient, do.ddl.OwnerManager().IsOwner) + if !do.shouldStartTTLJobManager() { + logutil.BgLogger().Info("skip starting ttl job manager for external workload role", + zap.String("role", string(do.extWorkloadMgr.Role()))) + return + } + ttlJobManager := ttlworker.NewJobManager(do.ddl.GetID(), do.advancedSysSessionPool, do.store, do.etcdClient, do.ddl.OwnerManager().IsOwner, do.extWorkloadMgr) do.ttlJobManager.Store(ttlJobManager) ttlJobManager.Start() } +func (do *Domain) shouldStartTTLJobManager() bool { + if !extworkload.IsEnabled(do.extWorkloadMgr) { + return true + } + return extworkload.IsTTLTaskWorker(do.extWorkloadMgr) +} + // TTLJobManager returns the ttl job manager on this domain func (do *Domain) TTLJobManager() *ttlworker.JobManager { return do.ttlJobManager.Load() diff --git a/pkg/domain/domain_sysvars.go b/pkg/domain/domain_sysvars.go index b03911f025ba6..3b8926d6ab326 100644 --- a/pkg/domain/domain_sysvars.go +++ b/pkg/domain/domain_sysvars.go @@ -19,6 +19,7 @@ import ( "strconv" "time" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/sessionctx/vardef" "github.com/pingcap/tidb/pkg/sessionctx/variable" "github.com/tikv/client-go/v2/tikv" @@ -48,6 +49,7 @@ func (do *Domain) initDomainSysVars() { variable.ChangeSchemaCacheSize = do.isSyncer.ChangeSchemaCacheSize variable.ChangePDMetadataCircuitBreakerErrorRateThresholdRatio = changePDMetadataCircuitBreakerErrorRateThresholdRatio + variable.UpdateExternalWorkloadTTLJobEnable = do.updateExternalWorkloadTTLJobEnable } // setStatsCacheCapacity sets statsCache cap @@ -124,6 +126,13 @@ func (*Domain) setGlobalResourceControl(enable bool) { } } +func (do *Domain) updateExternalWorkloadTTLJobEnable(ctx context.Context, enable bool) error { + if !extworkload.IsMaster(do.extWorkloadMgr) { + return nil + } + return do.extWorkloadMgr.UpdateTTLJobEnable(ctx, enable) +} + func (do *Domain) setLowResolutionTSOUpdateInterval(interval time.Duration) error { return do.store.GetOracle().SetLowResolutionTimestampUpdateInterval(interval) } diff --git a/pkg/domain/domain_test.go b/pkg/domain/domain_test.go index 30b92409f1be9..7077c713cd95c 100644 --- a/pkg/domain/domain_test.go +++ b/pkg/domain/domain_test.go @@ -27,7 +27,9 @@ import ( "github.com/ngaut/pools" "github.com/pingcap/errors" "github.com/pingcap/failpoint" + "github.com/pingcap/kvproto/pkg/keyspacepb" "github.com/pingcap/kvproto/pkg/metapb" + "github.com/pingcap/tidb/pkg/config" "github.com/pingcap/tidb/pkg/ddl" "github.com/pingcap/tidb/pkg/domain/infosync" "github.com/pingcap/tidb/pkg/domain/serverinfo" @@ -47,6 +49,49 @@ import ( "go.etcd.io/etcd/tests/v3/integration" ) +type fakeExternalWorkloadManager struct { + role config.ExternalWorkloadRole + updatedValue *bool +} + +func (m *fakeExternalWorkloadManager) Close() error { return nil } +func (m *fakeExternalWorkloadManager) Role() config.ExternalWorkloadRole { + return m.role +} +func (*fakeExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil } +func (*fakeExternalWorkloadManager) InitializeGCV2(context.Context) error { + return nil +} +func (*fakeExternalWorkloadManager) AbortGCV2(context.Context) error { return nil } +func (*fakeExternalWorkloadManager) RegisterGCV2(context.Context, uint64, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) UpdateGCLifeTime(context.Context, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RegisterTTLTask(context.Context, int64, bool) error { + return nil +} +func (*fakeExternalWorkloadManager) DeleteTTLTableInfo(context.Context, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleTTLTask(context.Context, uint64) error { + return nil +} +func (m *fakeExternalWorkloadManager) UpdateTTLJobEnable(_ context.Context, ttlJobEnable bool) error { + m.updatedValue = &ttlJobEnable + return nil +} +func (*fakeExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error { + return nil +} + func TestInfo(t *testing.T) { t.Skip("TestInfo will hang currently, it should be fixed later") @@ -216,6 +261,31 @@ func TestStatWorkRecoverFromPanic(t *testing.T) { require.True(t, isClose) } +func TestUpdateExternalWorkloadTTLJobEnableOnlyFromMaster(t *testing.T) { + dom := NewMockDomain() + master := &fakeExternalWorkloadManager{role: config.RoleMaster} + dom.SetExternalWorkloadManager(master) + require.NoError(t, dom.updateExternalWorkloadTTLJobEnable(context.Background(), false)) + require.NotNil(t, master.updatedValue) + require.False(t, *master.updatedValue) + + ttlWorker := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} + dom.SetExternalWorkloadManager(ttlWorker) + require.NoError(t, dom.updateExternalWorkloadTTLJobEnable(context.Background(), true)) + require.Nil(t, ttlWorker.updatedValue) +} + +func TestShouldStartTTLJobManagerWithExternalWorkloadRole(t *testing.T) { + dom := NewMockDomain() + require.True(t, dom.shouldStartTTLJobManager()) + + dom.SetExternalWorkloadManager(&fakeExternalWorkloadManager{role: config.RoleMaster}) + require.False(t, dom.shouldStartTTLJobManager()) + + dom.SetExternalWorkloadManager(&fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker}) + require.True(t, dom.shouldStartTTLJobManager()) +} + // ETCD use ip:port as unix socket address, however this address is invalid on windows. // We have to skip some of the test in such case. // https://github.com/etcd-io/etcd/blob/f0faa5501d936cd8c9f561bb9d1baca70eb67ab1/pkg/types/urls.go#L42 diff --git a/pkg/session/BUILD.bazel b/pkg/session/BUILD.bazel index 278f832530b28..f7d98732abf0b 100644 --- a/pkg/session/BUILD.bazel +++ b/pkg/session/BUILD.bazel @@ -47,6 +47,7 @@ go_library( "//pkg/expression/sessionexpr", "//pkg/extension", "//pkg/extension/extensionimpl", + "//pkg/extworkload", "//pkg/infoschema", "//pkg/infoschema/context", "//pkg/infoschema/issyncer", diff --git a/pkg/session/bootstrap_test.go b/pkg/session/bootstrap_test.go index f58c52fd74359..e46f179f69ef4 100644 --- a/pkg/session/bootstrap_test.go +++ b/pkg/session/bootstrap_test.go @@ -1555,7 +1555,7 @@ func makeStore(t *testing.T, keyspaceMeta *keyspacepb.KeyspaceMeta, isHasPrefix t.Cleanup(func() { ddl.CloseOwnerManager(mockStore) }) - dom, err := domap.getWithEtcdClient(mockStore, etcdClient, nil) + dom, err := domap.getWithEtcdClient(mockStore, etcdClient, nil, domainCreateOptions{}) require.NoError(t, err) defer dom.Close() diff --git a/pkg/session/session.go b/pkg/session/session.go index cd4c13a03ba79..69fa5ad694b91 100644 --- a/pkg/session/session.go +++ b/pkg/session/session.go @@ -66,6 +66,7 @@ import ( "github.com/pingcap/tidb/pkg/expression/sessionexpr" "github.com/pingcap/tidb/pkg/extension" "github.com/pingcap/tidb/pkg/extension/extensionimpl" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" infoschemactx "github.com/pingcap/tidb/pkg/infoschema/context" "github.com/pingcap/tidb/pkg/infoschema/issyncer" @@ -4293,12 +4294,17 @@ func InitMDLVariable(store kv.Storage) error { // BootstrapSession bootstrap session and domain. func BootstrapSession(store kv.Storage) (*domain.Domain, error) { - return bootstrapSessionImpl(context.Background(), store, createSessions) + return BootstrapSessionWithExternalWorkloadManager(store, nil) +} + +// BootstrapSessionWithExternalWorkloadManager bootstraps session and domain with an external workload manager. +func BootstrapSessionWithExternalWorkloadManager(store kv.Storage, manager extworkload.Manager) (*domain.Domain, error) { + return bootstrapSessionImpl(context.Background(), store, createSessions, manager) } // BootstrapSession4DistExecution bootstrap session and dom for Distributed execution test, only for unit testing. func BootstrapSession4DistExecution(store kv.Storage) (*domain.Domain, error) { - return bootstrapSessionImpl(context.Background(), store, createSessions4DistExecution) + return bootstrapSessionImpl(context.Background(), store, createSessions4DistExecution, nil) } // bootstrapSessionImpl bootstraps session and domain. @@ -4313,7 +4319,7 @@ func BootstrapSession4DistExecution(store kv.Storage) (*domain.Domain, error) { // - initialization global variables from system table that's required to use sessionCtx, // such as system time zone // - start domain and other routines. -func bootstrapSessionImpl(ctx context.Context, store kv.Storage, createSessionsImpl func(store kv.Storage, cnt int) ([]*session, error)) (*domain.Domain, error) { +func bootstrapSessionImpl(ctx context.Context, store kv.Storage, createSessionsImpl func(store kv.Storage, cnt int) ([]*session, error), extWorkloadMgr extworkload.Manager) (*domain.Domain, error) { ver := getStoreBootstrapVersionWithCache(store) if kv.IsUserKS(store) { targetVer := currentBootstrapVersion @@ -4360,7 +4366,6 @@ func bootstrapSessionImpl(ctx context.Context, store kv.Storage, createSessionsI return nil, err } } - // initiate disttask framework components which need a store scheduler.RegisterSchedulerFactory( proto.ImportInto, @@ -4384,6 +4389,11 @@ func bootstrapSessionImpl(ctx context.Context, store kv.Storage, createSessionsI concurrency = 0 } + if extWorkloadMgr != nil { + if _, err := domap.getWithEtcdClient(store, nil, nil, domainCreateOptions{extWorkloadMgr: extWorkloadMgr}); err != nil { + return nil, err + } + } ses, err := createSessionsImpl(store, 10) if err != nil { return nil, err diff --git a/pkg/session/session_test.go b/pkg/session/session_test.go index f5b8d08ee8900..9d3075311d89a 100644 --- a/pkg/session/session_test.go +++ b/pkg/session/session_test.go @@ -111,7 +111,7 @@ func TestBootstrapSessionImplUserKSVersionGuard(t *testing.T) { defer func() { panicVal = recover() }() - return bootstrapSessionImpl(context.Background(), userStore, createSessionStub) + return bootstrapSessionImpl(context.Background(), userStore, createSessionStub, nil) }() require.NotNil(t, panicVal) panicMsg := fmt.Sprint(panicVal) diff --git a/pkg/session/tidb.go b/pkg/session/tidb.go index 47c231e6f8c05..2870a702df979 100644 --- a/pkg/session/tidb.go +++ b/pkg/session/tidb.go @@ -31,6 +31,7 @@ import ( "github.com/pingcap/tidb/pkg/domain" "github.com/pingcap/tidb/pkg/errno" "github.com/pingcap/tidb/pkg/executor" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" "github.com/pingcap/tidb/pkg/infoschema/issyncer" "github.com/pingcap/tidb/pkg/infoschema/validatorapi" @@ -62,11 +63,15 @@ type domainMap struct { domains map[string]*domain.Domain } +type domainCreateOptions struct { + extWorkloadMgr extworkload.Manager +} + // Get or create the domain for store. // TODO decouple domain create from it, it's more clear to create domain explicitly // before any usage of it. func (dm *domainMap) Get(store kv.Storage) (d *domain.Domain, err error) { - return dm.getWithEtcdClient(store, nil, nil) + return dm.getWithEtcdClient(store, nil, nil, domainCreateOptions{}) } // GetOrCreateWithEtcdClient gets or creates the domain for store with etcd client. @@ -74,10 +79,10 @@ func (dm *domainMap) Get(store kv.Storage) (d *domain.Domain, err error) { // Caveat: If there is already a domain opened with your `store`, the filter passed in will be ignored and // the actual schema filter of the returned `Domain` is the one when the domain were created. func (dm *domainMap) GetOrCreateWithFilter(store kv.Storage, filter issyncer.Filter) (d *domain.Domain, err error) { - return dm.getWithEtcdClient(store, nil, filter) + return dm.getWithEtcdClient(store, nil, filter, domainCreateOptions{}) } -func (dm *domainMap) getWithEtcdClient(store kv.Storage, etcdClient *clientv3.Client, schemaFilter issyncer.Filter) (d *domain.Domain, err error) { +func (dm *domainMap) getWithEtcdClient(store kv.Storage, etcdClient *clientv3.Client, schemaFilter issyncer.Filter, opts domainCreateOptions) (d *domain.Domain, err error) { dm.mu.Lock() defer dm.mu.Unlock() @@ -113,6 +118,7 @@ func (dm *domainMap) getWithEtcdClient(store kv.Storage, etcdClient *clientv3.Cl etcdClient, schemaFilter, ) + d.SetExternalWorkloadManager(opts.extWorkloadMgr) var ddlInjector func(ddl.DDL, ddl.Executor, *infoschema.InfoCache) *schematracker.Checker if injector, ok := store.(schematracker.StorageDDLInjector); ok { diff --git a/pkg/sessionctx/variable/sysvar.go b/pkg/sessionctx/variable/sysvar.go index 714884a49b360..1ee5b3bc2d31f 100644 --- a/pkg/sessionctx/variable/sysvar.go +++ b/pkg/sessionctx/variable/sysvar.go @@ -3216,7 +3216,11 @@ var defaultSysVars = []*SysVar{ return BoolToOnOff(vardef.IgnoreInlistPlanDigest.Load()), nil }}, {Scope: vardef.ScopeGlobal, Name: vardef.TiDBTTLJobEnable, Value: BoolToOnOff(vardef.DefTiDBTTLJobEnable), Type: vardef.TypeBool, SetGlobal: func(ctx context.Context, vars *SessionVars, s string) error { - vardef.EnableTTLJob.Store(TiDBOptOn(s)) + enable := TiDBOptOn(s) + vardef.EnableTTLJob.Store(enable) + if UpdateExternalWorkloadTTLJobEnable != nil { + return UpdateExternalWorkloadTTLJobEnable(ctx, enable) + } return nil }, GetGlobal: func(ctx context.Context, vars *SessionVars) (string, error) { return BoolToOnOff(vardef.EnableTTLJob.Load()), nil diff --git a/pkg/sessionctx/variable/tidb_vars.go b/pkg/sessionctx/variable/tidb_vars.go index f0fc36d7a5f2f..30576c2c4373b 100644 --- a/pkg/sessionctx/variable/tidb_vars.go +++ b/pkg/sessionctx/variable/tidb_vars.go @@ -44,6 +44,8 @@ var ( GetExternalTimestamp func(ctx context.Context) (uint64, error) // SetGlobalResourceControl is the func registered by domain to set cluster resource control. SetGlobalResourceControl atomic.Pointer[func(bool)] + // UpdateExternalWorkloadTTLJobEnable is the func registered by domain to update external TTL worker scheduling. + UpdateExternalWorkloadTTLJobEnable func(ctx context.Context, enable bool) error // ValidateCloudStorageURI validates the cloud storage URI. ValidateCloudStorageURI func(ctx context.Context, uri string) error // SetLowResolutionTSOUpdateInterval is the func registered by domain to set slow resolution tso update interval. diff --git a/pkg/ttl/ttlworker/BUILD.bazel b/pkg/ttl/ttlworker/BUILD.bazel index 9348ad99a23f9..3039d15927215 100644 --- a/pkg/ttl/ttlworker/BUILD.bazel +++ b/pkg/ttl/ttlworker/BUILD.bazel @@ -17,11 +17,13 @@ go_library( importpath = "github.com/pingcap/tidb/pkg/ttl/ttlworker", visibility = ["//visibility:public"], deps = [ + "//pkg/extworkload", "//pkg/infoschema", "//pkg/infoschema/context", "//pkg/kv", "//pkg/meta/model", "//pkg/metrics", + "//pkg/owner", "//pkg/parser/ast", "//pkg/parser/terror", "//pkg/session/syssession", @@ -74,6 +76,7 @@ go_test( race = "on", shard_count = 50, deps = [ + "//pkg/config", "//pkg/domain", "//pkg/infoschema", "//pkg/infoschema/context", @@ -108,6 +111,7 @@ go_test( "@com_github_google_uuid//:uuid", "@com_github_pingcap_errors//:errors", "@com_github_pingcap_failpoint//:failpoint", + "@com_github_pingcap_kvproto//pkg/keyspacepb", "@com_github_prometheus_client_golang//prometheus", "@com_github_prometheus_client_model//go", "@com_github_stretchr_testify//assert", diff --git a/pkg/ttl/ttlworker/job_manager.go b/pkg/ttl/ttlworker/job_manager.go index 678e5f553843e..02cbb1807dffd 100644 --- a/pkg/ttl/ttlworker/job_manager.go +++ b/pkg/ttl/ttlworker/job_manager.go @@ -23,9 +23,11 @@ import ( "time" "github.com/pingcap/errors" + "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/infoschema" infoschemacontext "github.com/pingcap/tidb/pkg/infoschema/context" "github.com/pingcap/tidb/pkg/kv" + "github.com/pingcap/tidb/pkg/owner" "github.com/pingcap/tidb/pkg/parser/terror" "github.com/pingcap/tidb/pkg/session/syssession" "github.com/pingcap/tidb/pkg/sessionctx/vardef" @@ -45,6 +47,11 @@ import ( const scanTaskNotificationType string = "scan" +const ( + ttlJobManagerLeaderPath = "/tidb/ttl_job_manager/leader" + ttlJobManagerPrompt = "ttl_job_manager" +) + const insertNewTableIntoStatusTemplate = "INSERT INTO mysql.tidb_ttl_table_status (table_id,parent_table_id) VALUES (%?, %?)" const setTableStatusOwnerTemplate = `UPDATE mysql.tidb_ttl_table_status SET current_job_id = %?, @@ -126,14 +133,19 @@ type JobManager struct { lastReportDelayMetricsTime time.Time leaderFunc func() bool + ownerManager owner.Manager + extWorkload extworkload.Manager } // NewJobManager creates a new ttl job manager -func NewJobManager(id string, sessPool syssession.Pool, store kv.Storage, etcdCli *clientv3.Client, leaderFunc func() bool) (manager *JobManager) { +func NewJobManager(id string, sessPool syssession.Pool, store kv.Storage, etcdCli *clientv3.Client, leaderFunc func() bool, extWorkloadMgr ...extworkload.Manager) (manager *JobManager) { manager = &JobManager{} manager.id = id manager.store = store manager.sessPool = sessPool + if len(extWorkloadMgr) > 0 { + manager.extWorkload = extWorkloadMgr[0] + } manager.init(manager.jobLoop) manager.ctx = logutil.WithKeyValue(manager.ctx, "ttl-worker", "job-manager") @@ -156,15 +168,34 @@ func NewJobManager(id string, sessPool syssession.Pool, store kv.Storage, etcdCl manager.taskManager = newTaskManager(manager.ctx, sessPool, manager.infoSchemaCache, id, store) manager.leaderFunc = leaderFunc + if extworkload.IsTTLTaskWorker(manager.extWorkload) && etcdCli != nil && !intest.InTest { + manager.ownerManager = owner.NewOwnerManager(context.Background(), etcdCli, ttlJobManagerPrompt, id, ttlJobManagerLeaderPath) + manager.ownerManager.SetListener(&ttlJobManagerOwnerListener{}) + if err := manager.ownerManager.CampaignOwner(5); err != nil { + logutil.BgLogger().Error("failed to campaign ttl job manager owner", zap.Error(err)) + } + manager.leaderFunc = manager.ownerManager.IsOwner + } return } +type ttlJobManagerOwnerListener struct{} + +func (*ttlJobManagerOwnerListener) OnRetireOwner() {} + +func (*ttlJobManagerOwnerListener) OnBecomeOwner() { + logutil.BgLogger().Info("leader change of TTL job manager service, this node become owner") +} + func (m *JobManager) isLeader() bool { return m.leaderFunc != nil && m.leaderFunc() } func (m *JobManager) jobLoop() error { defer func() { + if m.ownerManager != nil { + m.ownerManager.Close() + } logutil.Logger(m.ctx).Info("ttlJobManager loop exited.") }() return withSession(m.sessPool, m.jobLoopWithSession) @@ -570,6 +601,9 @@ func (m *JobManager) findAllTasksForJob(se session.Session, jobID string) ([]*ca } func (m *JobManager) checkFinishedJob(se session.Session) { + runningJobsCount := len(m.runningJobs) + totalFinishedJobs := 0 + maxJobCreateTime := uint64(0) // reverse iteration so that we could remove the job safely in the loop for i := len(m.runningJobs) - 1; i >= 0; i-- { job := m.runningJobs[i] @@ -603,6 +637,17 @@ func (m *JobManager) checkFinishedJob(se session.Session) { continue } m.removeJob(job) + totalFinishedJobs++ + if createTime := uint64(job.createTime.Unix()); maxJobCreateTime < createTime { + maxJobCreateTime = createTime + } + } + } + if runningJobsCount > 0 && totalFinishedJobs == runningJobsCount && extworkload.IsTTLTaskWorker(m.extWorkload) { + if err := m.extWorkload.RecycleTTLTask(m.ctx, maxJobCreateTime); err != nil { + logutil.Logger(m.ctx).Warn("failed to recycle TTL task from external workload controller", + zap.Uint64("completedJobCreateTime", maxJobCreateTime), + zap.Error(err)) } } } diff --git a/pkg/ttl/ttlworker/job_manager_test.go b/pkg/ttl/ttlworker/job_manager_test.go index 0f99943a03247..54696cea78170 100644 --- a/pkg/ttl/ttlworker/job_manager_test.go +++ b/pkg/ttl/ttlworker/job_manager_test.go @@ -21,6 +21,8 @@ import ( "time" "github.com/pingcap/errors" + "github.com/pingcap/kvproto/pkg/keyspacepb" + "github.com/pingcap/tidb/pkg/config" "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/parser/ast" "github.com/pingcap/tidb/pkg/parser/mysql" @@ -37,6 +39,49 @@ import ( "github.com/tikv/client-go/v2/tikvrpc" ) +type fakeExternalWorkloadManager struct { + role config.ExternalWorkloadRole + recycledCreateTS uint64 +} + +func (m *fakeExternalWorkloadManager) Close() error { return nil } +func (m *fakeExternalWorkloadManager) Role() config.ExternalWorkloadRole { + return m.role +} +func (*fakeExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil } +func (*fakeExternalWorkloadManager) InitializeGCV2(context.Context) error { + return nil +} +func (*fakeExternalWorkloadManager) AbortGCV2(context.Context) error { return nil } +func (*fakeExternalWorkloadManager) RegisterGCV2(context.Context, uint64, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) UpdateGCLifeTime(context.Context, int64) error { + return nil +} +func (*fakeExternalWorkloadManager) RegisterTTLTask(context.Context, int64, bool) error { + return nil +} +func (*fakeExternalWorkloadManager) DeleteTTLTableInfo(context.Context, int64) error { + return nil +} +func (m *fakeExternalWorkloadManager) RecycleTTLTask(_ context.Context, completedJobCreateTime uint64) error { + m.recycledCreateTS = completedJobCreateTime + return nil +} +func (*fakeExternalWorkloadManager) UpdateTTLJobEnable(context.Context, bool) error { + return nil +} +func (*fakeExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error { + return nil +} +func (*fakeExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error { + return nil +} + func newTTLTableStatusRows(status ...*cache.TableStatus) []chunk.Row { c := chunk.NewChunkWithCapacity([]*types.FieldType{ types.NewFieldType(mysql.TypeLonglong), // table_id @@ -242,6 +287,57 @@ func (j *ttlJob) ID() string { return j.id } +func TestCheckFinishedJobRecyclesExternalTTLTask(t *testing.T) { + createTime := time.Unix(1234, 0) + externalMgr := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} + m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) + m.runningJobs = []*ttlJob{ + { + id: "job1", + ownerID: "test-id", + createTime: createTime, + tableID: 1, + status: cache.JobStatusRunning, + }, + } + + se := newMockSession(t) + sqlCounter := 0 + se.executeSQL = func(_ context.Context, sql string, args ...any) ([]chunk.Row, error) { + sqlCounter++ + if sqlCounter == 1 { + expectedSQL, expectedArgs := cache.SelectFromTTLTaskWithJobID("job1") + require.Equal(t, expectedSQL, sql) + require.Equal(t, expectedArgs, args) + } + return nil, nil + } + + m.CheckFinishedJob(se) + require.Empty(t, m.runningJobs) + require.Equal(t, uint64(createTime.Unix()), externalMgr.recycledCreateTS) + require.Equal(t, 4, sqlCounter) +} + +func TestCheckFinishedJobDoesNotRecycleExternalTTLTaskFromMaster(t *testing.T) { + externalMgr := &fakeExternalWorkloadManager{role: config.RoleMaster} + m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) + m.runningJobs = []*ttlJob{ + { + id: "job1", + ownerID: "test-id", + createTime: time.Unix(1234, 0), + tableID: 1, + status: cache.JobStatusRunning, + }, + } + se := newMockSession(t) + + m.CheckFinishedJob(se) + require.Empty(t, m.runningJobs) + require.Zero(t, externalMgr.recycledCreateTS) +} + func TestReadyForLockHBTimeoutJobTables(t *testing.T) { tbl := newMockTTLTbl(t, "t1") m := NewJobManager("test-id", nil, nil, nil, nil) From af54562844fc036224f50bb806d2d933406d9cc4 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Sun, 5 Jul 2026 23:06:22 +0800 Subject: [PATCH 02/16] fix log Signed-off-by: ystaticy --- pkg/ttl/ttlworker/job_manager.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/ttl/ttlworker/job_manager.go b/pkg/ttl/ttlworker/job_manager.go index 02cbb1807dffd..23dbedb64eb54 100644 --- a/pkg/ttl/ttlworker/job_manager.go +++ b/pkg/ttl/ttlworker/job_manager.go @@ -172,7 +172,8 @@ func NewJobManager(id string, sessPool syssession.Pool, store kv.Storage, etcdCl manager.ownerManager = owner.NewOwnerManager(context.Background(), etcdCli, ttlJobManagerPrompt, id, ttlJobManagerLeaderPath) manager.ownerManager.SetListener(&ttlJobManagerOwnerListener{}) if err := manager.ownerManager.CampaignOwner(5); err != nil { - logutil.BgLogger().Error("failed to campaign ttl job manager owner", zap.Error(err)) + logutil.BgLogger().Error("failed to campaign ttl job manager owner", + zap.Error(err)) } manager.leaderFunc = manager.ownerManager.IsOwner } From f652a113b2f00c4eaf2ea425093a1d7620e48ffd Mon Sep 17 00:00:00 2001 From: ystaticy Date: Tue, 28 Jul 2026 15:55:23 +0800 Subject: [PATCH 03/16] fix comments Signed-off-by: ystaticy --- pkg/ddl/create_table.go | 4 +--- pkg/ddl/ttl.go | 25 ++++++++++++++++++----- pkg/ddl/ttl_test.go | 20 ++++++++++++++++++ pkg/session/session.go | 12 +++++++---- pkg/session/session_nextgen_test.go | 22 +++++++++++++++++++- pkg/sessionctx/variable/sysvar.go | 6 ++++-- pkg/sessionctx/variable/sysvar_test.go | 28 ++++++++++++++++++++++++++ 7 files changed, 102 insertions(+), 15 deletions(-) diff --git a/pkg/ddl/create_table.go b/pkg/ddl/create_table.go index ae73381ffbec5..1b6e555447fbe 100644 --- a/pkg/ddl/create_table.go +++ b/pkg/ddl/create_table.go @@ -249,9 +249,7 @@ func (w *worker) onCreateTable(jobCtx *jobContext, job *model.Job) (ver int64, _ return ver, errors.Trace(err) } - if err := w.registerTTLTableToExternalWorkload(jobCtx.ctx, tbInfo); err != nil { - return ver, errors.Trace(err) - } + w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo) // Finish this job. job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo) diff --git a/pkg/ddl/ttl.go b/pkg/ddl/ttl.go index f3adfee3c51ca..ac538f6615285 100644 --- a/pkg/ddl/ttl.go +++ b/pkg/ddl/ttl.go @@ -96,11 +96,7 @@ func onTTLInfoChange(jobCtx *jobContext, job *model.Job) (ver int64, err error) return ver, errors.Trace(err) } if jobCtx.oldDDLCtx != nil { - if err := jobCtx.oldDDLCtx.registerTTLTableToExternalWorkload(jobCtx.ctx, tblInfo); err != nil { - logutil.DDLLogger().Warn("failed to register TTL table to external workload controller", - zap.Int64("tableID", tblInfo.ID), - zap.Error(err)) - } + jobCtx.oldDDLCtx.syncTTLTableToExternalWorkload(jobCtx.ctx, tblInfo) } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) return ver, nil @@ -122,6 +118,25 @@ func (dc *ddlCtx) registerTTLTableToExternalWorkload(ctx context.Context, tblInf return manager.RegisterTTLTask(ctx, tblInfo.ID, vardef.EnableTTLJob.Load()) } +func (dc *ddlCtx) tryRegisterTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) { + if err := dc.registerTTLTableToExternalWorkload(ctx, tblInfo); err != nil { + logutil.DDLLogger().Warn("failed to register TTL table to external workload controller", + zap.Int64("tableID", tblInfo.ID), + zap.Error(err)) + } +} + +func (dc *ddlCtx) syncTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) { + if tblInfo == nil { + return + } + if tblInfo.TTLInfo == nil || !tblInfo.TTLInfo.Enable { + dc.deleteTTLTableFromExternalWorkload(ctx, tblInfo.ID) + return + } + dc.tryRegisterTTLTableToExternalWorkload(ctx, tblInfo) +} + func (dc *ddlCtx) deleteTTLTableFromExternalWorkload(ctx context.Context, tableID int64) { manager, ok := dc.externalWorkloadMaster() if !ok { diff --git a/pkg/ddl/ttl_test.go b/pkg/ddl/ttl_test.go index 93f3708da4605..112d37d7eabaf 100644 --- a/pkg/ddl/ttl_test.go +++ b/pkg/ddl/ttl_test.go @@ -257,3 +257,23 @@ func TestExternalWorkloadTTLTableRegisterReturnsError(t *testing.T) { }) require.ErrorIs(t, err, boom) } + +func TestExternalWorkloadTTLTableTryRegisterSwallowsError(t *testing.T) { + manager := &fakeExternalWorkloadManager{role: config.RoleMaster, registerErr: errors.New("boom")} + dc := &ddlCtx{extWorkload: manager} + dc.tryRegisterTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ + ID: 123, + TTLInfo: &model.TTLInfo{Enable: true}, + }) + require.Equal(t, int64(123), manager.registeredTable) +} + +func TestExternalWorkloadTTLTableSyncDeletesDisabledTTL(t *testing.T) { + manager := &fakeExternalWorkloadManager{role: config.RoleMaster} + dc := &ddlCtx{extWorkload: manager} + dc.syncTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ + ID: 123, + TTLInfo: &model.TTLInfo{Enable: false}, + }) + require.Equal(t, int64(123), manager.deletedTable) +} diff --git a/pkg/session/session.go b/pkg/session/session.go index a7c8f45db51d4..ddec599c3a90e 100644 --- a/pkg/session/session.go +++ b/pkg/session/session.go @@ -4449,7 +4449,7 @@ func bootstrapSessionImpl(ctx context.Context, store kv.Storage, createSessionsI return nil, err } if ver < currentBootstrapVersion { - err = runInBootstrapSession(store, ver) + err = runInBootstrapSession(store, ver, domainCreateOptions{extWorkloadMgr: extWorkloadMgr}) if err != nil { return nil, err } @@ -4704,7 +4704,7 @@ func getStartMode(ver int64) ddl.StartMode { // If no bootstrap and storage is remote, we must use a little lease time to // bootstrap quickly, after bootstrapped, we will reset the lease time. // TODO: Using a bootstrap tool for doing this may be better later. -func runInBootstrapSession(store kv.Storage, ver int64) error { +func runInBootstrapSession(store kv.Storage, ver int64, opts domainCreateOptions) error { startMode := getStartMode(ver) startTime := time.Now() defer func() { @@ -4751,7 +4751,7 @@ func runInBootstrapSession(store kv.Storage, ver int64) error { } } } - s, err := createSession(store) + s, err := createSessionWithDomainOptions(store, opts) if err != nil { // Bootstrap fail will cause program exit. logutil.BgLogger().Fatal("createSession error", zap.Error(err)) @@ -4823,7 +4823,11 @@ func createSessionsImpl(store kv.Storage, cnt int) ([]*session, error) { // This means the min ts reporter is not aware of it and may report a wrong min start ts. // In most cases you should use a session pool in domain instead. func createSession(store kv.Storage) (*session, error) { - dom, err := domap.Get(store) + return createSessionWithDomainOptions(store, domainCreateOptions{}) +} + +func createSessionWithDomainOptions(store kv.Storage, opts domainCreateOptions) (*session, error) { + dom, err := domap.getWithEtcdClient(store, nil, nil, opts) if err != nil { return nil, err } diff --git a/pkg/session/session_nextgen_test.go b/pkg/session/session_nextgen_test.go index 5a5beb439b596..e088d44200020 100644 --- a/pkg/session/session_nextgen_test.go +++ b/pkg/session/session_nextgen_test.go @@ -32,6 +32,10 @@ type upgradeGCV2Manager struct { abortCount int } +type bootstrapExternalWorkloadManager struct { + extworkload.Manager +} + func (*upgradeGCV2Manager) Role() config.ExternalWorkloadRole { return config.RoleGCV2Worker } @@ -74,7 +78,23 @@ func TestUpgradeGCV2AbortUsesPostLockBootstrapVersion(t *testing.T) { mgr := &upgradeGCV2Manager{} extworkload.SetManagerForStore(store, mgr) - runInBootstrapSession(store, currentBootstrapVersion-1) + runInBootstrapSession(store, currentBootstrapVersion-1, domainCreateOptions{}) require.Zero(t, mgr.abortCount) } + +func TestCreateSessionWithDomainOptionsAttachesExternalWorkloadManager(t *testing.T) { + store, dom := CreateStoreAndBootstrap(t) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + dom.Close() + domap.Delete(store) + + mgr := &bootstrapExternalWorkloadManager{} + newDom, err := domap.getWithEtcdClient(store, nil, nil, domainCreateOptions{extWorkloadMgr: mgr}) + require.NoError(t, err) + require.Same(t, mgr, newDom.ExternalWorkloadManager()) + + newDom.Close() +} diff --git a/pkg/sessionctx/variable/sysvar.go b/pkg/sessionctx/variable/sysvar.go index 4b6a9c0c32d86..a0883820a5679 100644 --- a/pkg/sessionctx/variable/sysvar.go +++ b/pkg/sessionctx/variable/sysvar.go @@ -3251,10 +3251,12 @@ var defaultSysVars = []*SysVar{ }}, {Scope: vardef.ScopeGlobal, Name: vardef.TiDBTTLJobEnable, Value: BoolToOnOff(vardef.DefTiDBTTLJobEnable), Type: vardef.TypeBool, SetGlobal: func(ctx context.Context, vars *SessionVars, s string) error { enable := TiDBOptOn(s) - vardef.EnableTTLJob.Store(enable) if UpdateExternalWorkloadTTLJobEnable != nil { - return UpdateExternalWorkloadTTLJobEnable(ctx, enable) + if err := UpdateExternalWorkloadTTLJobEnable(ctx, enable); err != nil { + return err + } } + vardef.EnableTTLJob.Store(enable) return nil }, GetGlobal: func(ctx context.Context, vars *SessionVars) (string, error) { return BoolToOnOff(vardef.EnableTTLJob.Load()), nil diff --git a/pkg/sessionctx/variable/sysvar_test.go b/pkg/sessionctx/variable/sysvar_test.go index b61b658538ba2..7504168cc918a 100644 --- a/pkg/sessionctx/variable/sysvar_test.go +++ b/pkg/sessionctx/variable/sysvar_test.go @@ -252,6 +252,34 @@ func TestTiFlashQuerySpillRatio(t *testing.T) { require.Equal(t, 0.75, vars.TiFlashQuerySpillRatio) } +func TestTiDBTTLJobEnableSetGlobalOnlyUpdatesLocalOnSuccess(t *testing.T) { + vars := NewSessionVars(nil) + sv := GetSysVar(vardef.TiDBTTLJobEnable) + require.NotNil(t, sv) + + originalEnable := vardef.EnableTTLJob.Load() + originalHook := UpdateExternalWorkloadTTLJobEnable + t.Cleanup(func() { + vardef.EnableTTLJob.Store(originalEnable) + UpdateExternalWorkloadTTLJobEnable = originalHook + }) + + vardef.EnableTTLJob.Store(false) + boom := fmt.Errorf("boom") + UpdateExternalWorkloadTTLJobEnable = func(context.Context, bool) error { + return boom + } + err := sv.SetGlobal(context.Background(), vars, vardef.On) + require.ErrorIs(t, err, boom) + require.False(t, vardef.EnableTTLJob.Load()) + + UpdateExternalWorkloadTTLJobEnable = func(context.Context, bool) error { + return nil + } + require.NoError(t, sv.SetGlobal(context.Background(), vars, vardef.On)) + require.True(t, vardef.EnableTTLJob.Load()) +} + func TestTiFlashHashJoinVersion(t *testing.T) { vars := NewSessionVars(nil) sv := GetSysVar(vardef.TiFlashHashJoinVersion) From f12b98295078cf00c0c50acbcd62aead4141a119 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Tue, 28 Jul 2026 16:18:03 +0800 Subject: [PATCH 04/16] ttl: register external workload for fk ttl tables --- pkg/ddl/BUILD.bazel | 1 + pkg/ddl/create_table.go | 1 + pkg/ddl/external_workload_ttl_test.go | 154 ++++++++++++++++++++++++++ 3 files changed, 156 insertions(+) create mode 100644 pkg/ddl/external_workload_ttl_test.go diff --git a/pkg/ddl/BUILD.bazel b/pkg/ddl/BUILD.bazel index 5fcbb43c3ee79..25cb25d32e669 100644 --- a/pkg/ddl/BUILD.bazel +++ b/pkg/ddl/BUILD.bazel @@ -272,6 +272,7 @@ go_test( "executor_nokit_test.go", "executor_test.go", "export_test.go", + "external_workload_ttl_test.go", "fail_test.go", "foreign_key_test.go", "index_change_test.go", diff --git a/pkg/ddl/create_table.go b/pkg/ddl/create_table.go index 1b6e555447fbe..7c302a7dc1e4b 100644 --- a/pkg/ddl/create_table.go +++ b/pkg/ddl/create_table.go @@ -295,6 +295,7 @@ func (w *worker) createTableWithForeignKeys(jobCtx *jobContext, job *model.Job, if err != nil { return ver, errors.Trace(err) } + w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo) job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo) return ver, nil diff --git a/pkg/ddl/external_workload_ttl_test.go b/pkg/ddl/external_workload_ttl_test.go new file mode 100644 index 0000000000000..16baf9e66db3f --- /dev/null +++ b/pkg/ddl/external_workload_ttl_test.go @@ -0,0 +1,154 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package ddl_test + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/pingcap/kvproto/pkg/keyspacepb" + "github.com/pingcap/tidb/pkg/config" + "github.com/pingcap/tidb/pkg/domain" + "github.com/pingcap/tidb/pkg/kv" + "github.com/pingcap/tidb/pkg/session" + "github.com/pingcap/tidb/pkg/sessionctx/vardef" + "github.com/pingcap/tidb/pkg/store/mockstore" + "github.com/pingcap/tidb/pkg/testkit" + "github.com/pingcap/tidb/pkg/testkit/external" + "github.com/stretchr/testify/require" +) + +type recordingExternalWorkloadManager struct { + role config.ExternalWorkloadRole + + mu sync.Mutex + registeredTables []int64 + deletedTables []int64 +} + +func (*recordingExternalWorkloadManager) Close() error { return nil } + +func (m *recordingExternalWorkloadManager) Role() config.ExternalWorkloadRole { + return m.role +} + +func (*recordingExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil } + +func (*recordingExternalWorkloadManager) InitializeGCV2(context.Context, time.Duration) error { + return nil +} + +func (*recordingExternalWorkloadManager) AbortGCV2(context.Context) error { return nil } + +func (*recordingExternalWorkloadManager) RegisterGCV2(context.Context, uint64, time.Duration) error { + return nil +} + +func (*recordingExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error { + return nil +} + +func (*recordingExternalWorkloadManager) UpdateGCLifeTime(context.Context, time.Duration) error { + return nil +} + +func (m *recordingExternalWorkloadManager) RegisterTTLTask(_ context.Context, tableID int64, _ bool) error { + m.mu.Lock() + defer m.mu.Unlock() + m.registeredTables = append(m.registeredTables, tableID) + return nil +} + +func (m *recordingExternalWorkloadManager) DeleteTTLTableInfo(_ context.Context, tableID int64) error { + m.mu.Lock() + defer m.mu.Unlock() + m.deletedTables = append(m.deletedTables, tableID) + return nil +} + +func (*recordingExternalWorkloadManager) RecycleTTLTask(context.Context, uint64) error { + return nil +} + +func (*recordingExternalWorkloadManager) UpdateTTLJobEnable(context.Context, bool) error { + return nil +} + +func (*recordingExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error { + return nil +} + +func (*recordingExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error { + return nil +} + +func (m *recordingExternalWorkloadManager) registeredTTLTables() []int64 { + m.mu.Lock() + defer m.mu.Unlock() + return append([]int64(nil), m.registeredTables...) +} + +func (m *recordingExternalWorkloadManager) deletedTTLTables() []int64 { + m.mu.Lock() + defer m.mu.Unlock() + return append([]int64(nil), m.deletedTables...) +} + +func createTTLExternalWorkloadTestKit(t *testing.T, mgr *recordingExternalWorkloadManager) (*testkit.TestKit, kv.Storage) { + store, err := mockstore.NewMockStore() + require.NoError(t, err) + + vardef.SetSchemaLease(500 * time.Millisecond) + session.DisableStats4Test() + domain.DisablePlanReplayerBackgroundJob4Test() + domain.DisableDumpHistoricalStats4Test() + + dom, err := session.BootstrapSessionWithExternalWorkloadManager(store, mgr) + require.NoError(t, err) + dom.SetStatsUpdating(true) + + t.Cleanup(func() { + dom.Close() + require.NoError(t, store.Close()) + }) + + tk := testkit.NewTestKit(t, store) + tk.MustExec("use test") + tk.MustExec("set @@global.tidb_enable_foreign_key=1") + tk.MustExec("set @@foreign_key_checks=1") + return tk, store +} + +func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { + t.Run("create table with foreign key registers ttl", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec("create table parent(id int primary key)") + tk.MustExec(`create table child( + id int primary key, + parent_id int, + created_at datetime, + index idx_parent(parent_id), + foreign key (parent_id) references parent(id) + ) TTL = created_at + interval 1 day`) + + childTbl := external.GetTableByName(t, tk, "test", "child") + require.Equal(t, []int64{childTbl.Meta().ID}, mgr.registeredTTLTables()) + require.Empty(t, mgr.deletedTTLTables()) + }) +} From aec986619d5f808536eef150440a5409b6f83e0a Mon Sep 17 00:00:00 2001 From: ystaticy Date: Tue, 28 Jul 2026 16:21:12 +0800 Subject: [PATCH 05/16] ttl: sync external metadata on drop and truncate --- pkg/ddl/external_workload_ttl_test.go | 33 +++++++++++++++++++++++++++ pkg/ddl/table.go | 11 +++++++++ 2 files changed, 44 insertions(+) diff --git a/pkg/ddl/external_workload_ttl_test.go b/pkg/ddl/external_workload_ttl_test.go index 16baf9e66db3f..b7ca82df83801 100644 --- a/pkg/ddl/external_workload_ttl_test.go +++ b/pkg/ddl/external_workload_ttl_test.go @@ -151,4 +151,37 @@ func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { require.Equal(t, []int64{childTbl.Meta().ID}, mgr.registeredTTLTables()) require.Empty(t, mgr.deletedTTLTables()) }) + + t.Run("drop table deletes ttl metadata", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + tbl := external.GetTableByName(t, tk, "test", "t") + tk.MustExec("drop table t") + + require.Equal(t, []int64{tbl.Meta().ID}, mgr.deletedTTLTables()) + }) + + t.Run("truncate table refreshes ttl metadata", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + oldTbl := external.GetTableByName(t, tk, "test", "t") + tk.MustExec("truncate table t") + newTbl := external.GetTableByName(t, tk, "test", "t") + + require.NotEqual(t, oldTbl.Meta().ID, newTbl.Meta().ID) + require.Equal(t, []int64{oldTbl.Meta().ID}, mgr.deletedTTLTables()) + require.Equal(t, []int64{oldTbl.Meta().ID, newTbl.Meta().ID}, mgr.registeredTTLTables()) + }) } diff --git a/pkg/ddl/table.go b/pkg/ddl/table.go index e11bcf3c44f11..1e0da1b67c5ab 100644 --- a/pkg/ddl/table.go +++ b/pkg/ddl/table.go @@ -140,6 +140,9 @@ func (w *worker) onDropTableOrView(jobCtx *jobContext, job *model.Job) (ver int6 if err := w.dropMaskingPoliciesOnTable(jobCtx, tblInfo.ID); err != nil { return ver, errors.Wrapf(err, "failed to drop masking policies on table %d", tblInfo.ID) } + if jobCtx.oldDDLCtx != nil && tblInfo.TTLInfo != nil { + jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID) + } // Finish this job. job.FinishTableJob(model.JobStateDone, model.StateNone, ver, tblInfo) @@ -626,6 +629,14 @@ func (w *worker) onTruncateTable(jobCtx *jobContext, job *model.Job) (ver int64, if err != nil { return ver, errors.Trace(err) } + if jobCtx.oldDDLCtx != nil { + if oldTblInfo.TTLInfo != nil { + jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, oldTblInfo.ID) + } + if tblInfo.TTLInfo != nil && tblInfo.TTLInfo.Enable { + jobCtx.oldDDLCtx.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tblInfo) + } + } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) // see truncateTableByReassignPartitionIDs for why they might change. From d992297487dcfb595f36d101e6923960f1fdf3b5 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Tue, 28 Jul 2026 16:24:27 +0800 Subject: [PATCH 06/16] ttl: allow ddl sync from ttl worker role --- pkg/ddl/external_workload_ttl_test.go | 16 ++++++++++++++++ pkg/ddl/ttl.go | 8 ++++---- pkg/ddl/ttl_test.go | 4 ++-- 3 files changed, 22 insertions(+), 6 deletions(-) diff --git a/pkg/ddl/external_workload_ttl_test.go b/pkg/ddl/external_workload_ttl_test.go index b7ca82df83801..493b258a03503 100644 --- a/pkg/ddl/external_workload_ttl_test.go +++ b/pkg/ddl/external_workload_ttl_test.go @@ -184,4 +184,20 @@ func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { require.Equal(t, []int64{oldTbl.Meta().ID}, mgr.deletedTTLTables()) require.Equal(t, []int64{oldTbl.Meta().ID, newTbl.Meta().ID}, mgr.registeredTTLTables()) }) + + t.Run("ddl syncs ttl metadata from ttl worker role", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleTTLTaskWorker} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + tbl := external.GetTableByName(t, tk, "test", "t") + require.Equal(t, []int64{tbl.Meta().ID}, mgr.registeredTTLTables()) + + tk.MustExec("drop table t") + require.Equal(t, []int64{tbl.Meta().ID}, mgr.deletedTTLTables()) + }) } diff --git a/pkg/ddl/ttl.go b/pkg/ddl/ttl.go index ac538f6615285..a5169bb010ab5 100644 --- a/pkg/ddl/ttl.go +++ b/pkg/ddl/ttl.go @@ -102,16 +102,16 @@ func onTTLInfoChange(jobCtx *jobContext, job *model.Job) (ver int64, err error) return ver, nil } -func (dc *ddlCtx) externalWorkloadMaster() (extworkload.Manager, bool) { +func (dc *ddlCtx) externalWorkloadManager() (extworkload.Manager, bool) { if dc == nil { return nil, false } manager := dc.extWorkload - return manager, extworkload.IsMaster(manager) + return manager, extworkload.IsEnabled(manager) } func (dc *ddlCtx) registerTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) error { - manager, ok := dc.externalWorkloadMaster() + manager, ok := dc.externalWorkloadManager() if !ok || tblInfo == nil || tblInfo.TTLInfo == nil || !tblInfo.TTLInfo.Enable { return nil } @@ -138,7 +138,7 @@ func (dc *ddlCtx) syncTTLTableToExternalWorkload(ctx context.Context, tblInfo *m } func (dc *ddlCtx) deleteTTLTableFromExternalWorkload(ctx context.Context, tableID int64) { - manager, ok := dc.externalWorkloadMaster() + manager, ok := dc.externalWorkloadManager() if !ok { return } diff --git a/pkg/ddl/ttl_test.go b/pkg/ddl/ttl_test.go index 112d37d7eabaf..712acbac57f15 100644 --- a/pkg/ddl/ttl_test.go +++ b/pkg/ddl/ttl_test.go @@ -233,8 +233,8 @@ func TestExternalWorkloadTTLTableReportsOnlyFromMaster(t *testing.T) { dc = &ddlCtx{extWorkload: ttlWorker} require.NoError(t, dc.registerTTLTableToExternalWorkload(context.Background(), tblInfo)) dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID) - require.Zero(t, ttlWorker.registeredTable) - require.Zero(t, ttlWorker.deletedTable) + require.Equal(t, int64(123), ttlWorker.registeredTable) + require.Equal(t, int64(123), ttlWorker.deletedTable) } func TestExternalWorkloadTTLTableRegisterSkipsDisabledTTL(t *testing.T) { From 6d3541525ca7fe64fb585933f7f223bfc0ac1301 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Tue, 28 Jul 2026 16:33:27 +0800 Subject: [PATCH 07/16] session: cover bootstrap external workload manager hookup --- pkg/session/session.go | 1 + pkg/session/session_nextgen_test.go | 82 ++++++++++++++++++++++++++++- 2 files changed, 81 insertions(+), 2 deletions(-) diff --git a/pkg/session/session.go b/pkg/session/session.go index ddec599c3a90e..601c9fe4e4755 100644 --- a/pkg/session/session.go +++ b/pkg/session/session.go @@ -4757,6 +4757,7 @@ func runInBootstrapSession(store kv.Storage, ver int64, opts domainCreateOptions logutil.BgLogger().Fatal("createSession error", zap.Error(err)) } dom := domain.GetDomain(s) + failpoint.InjectCall("checkBootstrapExternalWorkloadManager", dom) err = dom.Start(startMode) if err != nil { // Bootstrap fail will cause program exit. diff --git a/pkg/session/session_nextgen_test.go b/pkg/session/session_nextgen_test.go index e088d44200020..b234340f90d20 100644 --- a/pkg/session/session_nextgen_test.go +++ b/pkg/session/session_nextgen_test.go @@ -19,11 +19,16 @@ package session import ( "context" "testing" + "time" + "github.com/pingcap/kvproto/pkg/keyspacepb" "github.com/pingcap/tidb/pkg/config" "github.com/pingcap/tidb/pkg/config/deploymode" + "github.com/pingcap/tidb/pkg/domain" "github.com/pingcap/tidb/pkg/extworkload" "github.com/pingcap/tidb/pkg/sessionctx/variable" + "github.com/pingcap/tidb/pkg/store/mockstore" + "github.com/pingcap/tidb/pkg/testkit/testfailpoint" "github.com/stretchr/testify/require" ) @@ -33,7 +38,7 @@ type upgradeGCV2Manager struct { } type bootstrapExternalWorkloadManager struct { - extworkload.Manager + role config.ExternalWorkloadRole } func (*upgradeGCV2Manager) Role() config.ExternalWorkloadRole { @@ -45,6 +50,56 @@ func (m *upgradeGCV2Manager) AbortGCV2(context.Context) error { return nil } +func (*bootstrapExternalWorkloadManager) Close() error { return nil } + +func (m *bootstrapExternalWorkloadManager) Role() config.ExternalWorkloadRole { + return m.role +} + +func (*bootstrapExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil } + +func (*bootstrapExternalWorkloadManager) InitializeGCV2(context.Context, time.Duration) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) AbortGCV2(context.Context) error { return nil } + +func (*bootstrapExternalWorkloadManager) RegisterGCV2(context.Context, uint64, time.Duration) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) UpdateGCLifeTime(context.Context, time.Duration) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) RegisterTTLTask(context.Context, int64, bool) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) DeleteTTLTableInfo(context.Context, int64) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) RecycleTTLTask(context.Context, uint64) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) UpdateTTLJobEnable(context.Context, bool) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error { + return nil +} + +func (*bootstrapExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error { + return nil +} + func TestUsePipelinedDMLDisabledInStarter(t *testing.T) { originalMode := deploymode.Get() require.NoError(t, deploymode.Set(deploymode.Starter)) @@ -91,10 +146,33 @@ func TestCreateSessionWithDomainOptionsAttachesExternalWorkloadManager(t *testin dom.Close() domap.Delete(store) - mgr := &bootstrapExternalWorkloadManager{} + mgr := &bootstrapExternalWorkloadManager{role: config.RoleMaster} newDom, err := domap.getWithEtcdClient(store, nil, nil, domainCreateOptions{extWorkloadMgr: mgr}) require.NoError(t, err) require.Same(t, mgr, newDom.ExternalWorkloadManager()) newDom.Close() } + +func TestBootstrapSessionWithExternalWorkloadManagerAttachesBootstrapDomain(t *testing.T) { + store, err := mockstore.NewMockStore() + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, store.Close()) + }) + + mgr := &bootstrapExternalWorkloadManager{role: config.RoleMaster} + sawBootstrapDomain := false + testfailpoint.EnableCall(t, "github.com/pingcap/tidb/pkg/session/checkBootstrapExternalWorkloadManager", func(bootstrapDom *domain.Domain) { + require.Equal(t, extworkload.Manager(mgr), bootstrapDom.ExternalWorkloadManager()) + sawBootstrapDomain = true + }) + + newDom, err := BootstrapSessionWithExternalWorkloadManager(store, mgr) + require.NoError(t, err) + require.True(t, sawBootstrapDomain) + require.Equal(t, extworkload.Manager(mgr), newDom.ExternalWorkloadManager()) + + newDom.Close() + domap.Delete(store) +} From 566e04d8832eea216757e9a3d217560b963d283e Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 09:29:17 +0800 Subject: [PATCH 08/16] fix comments Signed-off-by: ystaticy --- pkg/domain/domain.go | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/pkg/domain/domain.go b/pkg/domain/domain.go index 220e068777bb2..5a08e35a5e565 100644 --- a/pkg/domain/domain.go +++ b/pkg/domain/domain.go @@ -2890,7 +2890,7 @@ func (do *Domain) serverIDKeeper() { // StartTTLJobManager creates and starts the ttl job manager func (do *Domain) StartTTLJobManager() { if !do.shouldStartTTLJobManager() { - logutil.BgLogger().Info("skip starting ttl job manager for external workload role", + logutil.BgLogger().Info("don't run ttl job manager", zap.String("role", string(do.extWorkloadMgr.Role()))) return } @@ -2900,10 +2900,13 @@ func (do *Domain) StartTTLJobManager() { } func (do *Domain) shouldStartTTLJobManager() bool { - if !extworkload.IsEnabled(do.extWorkloadMgr) { - return true + // Unit tests and non-worker deployments do not install an external workload + // manager, so they still start a local TTL job manager. Once external + // workload is enabled, only the dedicated TTL task worker should run TTL jobs. + if extworkload.IsEnabled(do.extWorkloadMgr) && !extworkload.IsTTLTaskWorker(do.extWorkloadMgr) { + return false } - return extworkload.IsTTLTaskWorker(do.extWorkloadMgr) + return true } // TTLJobManager returns the ttl job manager on this domain From 0c4cfd4f6029444d8406e26127633eb408d9408c Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 09:41:40 +0800 Subject: [PATCH 09/16] fix comments Signed-off-by: ystaticy --- pkg/session/session_nextgen_test.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/pkg/session/session_nextgen_test.go b/pkg/session/session_nextgen_test.go index b234340f90d20..a0d9f58c47886 100644 --- a/pkg/session/session_nextgen_test.go +++ b/pkg/session/session_nextgen_test.go @@ -147,10 +147,12 @@ func TestCreateSessionWithDomainOptionsAttachesExternalWorkloadManager(t *testin domap.Delete(store) mgr := &bootstrapExternalWorkloadManager{role: config.RoleMaster} - newDom, err := domap.getWithEtcdClient(store, nil, nil, domainCreateOptions{extWorkloadMgr: mgr}) + se, err := createSessionWithDomainOptions(store, domainCreateOptions{extWorkloadMgr: mgr}) require.NoError(t, err) + newDom := domain.GetDomain(se) require.Same(t, mgr, newDom.ExternalWorkloadManager()) + se.Close() newDom.Close() } From 6c4355b88c4bbd3cc6804adab0f6d0cd0c778172 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 09:52:57 +0800 Subject: [PATCH 10/16] fix comments Signed-off-by: ystaticy --- pkg/domain/domain.go | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pkg/domain/domain.go b/pkg/domain/domain.go index 5a08e35a5e565..48cf1350d1ec0 100644 --- a/pkg/domain/domain.go +++ b/pkg/domain/domain.go @@ -235,7 +235,10 @@ type Domain struct { // only used for nextgen crossKSSessMgr *crossks.Manager crossKSSessFactoryGetter func(string, validatorapi.Validator) pools.Factory - extWorkloadMgr extworkload.Manager + + // extWorkloadMgr coordinates background workloads such as TTL with the + // external workload controller in Starter deployments. + extWorkloadMgr extworkload.Manager } var _ sqlsvrapi.Server = (*Domain)(nil) From 0d2c0192927e73242c2a6b2a2a40cfe6c6aad3bf Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 10:02:51 +0800 Subject: [PATCH 11/16] fix comments Signed-off-by: ystaticy --- pkg/ttl/ttlworker/job_manager.go | 3 +- pkg/ttl/ttlworker/job_manager_test.go | 202 ++++++++++++++++++++++---- 2 files changed, 177 insertions(+), 28 deletions(-) diff --git a/pkg/ttl/ttlworker/job_manager.go b/pkg/ttl/ttlworker/job_manager.go index 23dbedb64eb54..e2386c8cf9f4a 100644 --- a/pkg/ttl/ttlworker/job_manager.go +++ b/pkg/ttl/ttlworker/job_manager.go @@ -602,7 +602,6 @@ func (m *JobManager) findAllTasksForJob(se session.Session, jobID string) ([]*ca } func (m *JobManager) checkFinishedJob(se session.Session) { - runningJobsCount := len(m.runningJobs) totalFinishedJobs := 0 maxJobCreateTime := uint64(0) // reverse iteration so that we could remove the job safely in the loop @@ -644,7 +643,7 @@ func (m *JobManager) checkFinishedJob(se session.Session) { } } } - if runningJobsCount > 0 && totalFinishedJobs == runningJobsCount && extworkload.IsTTLTaskWorker(m.extWorkload) { + if totalFinishedJobs > 0 && extworkload.IsTTLTaskWorker(m.extWorkload) { if err := m.extWorkload.RecycleTTLTask(m.ctx, maxJobCreateTime); err != nil { logutil.Logger(m.ctx).Warn("failed to recycle TTL task from external workload controller", zap.Uint64("completedJobCreateTime", maxJobCreateTime), diff --git a/pkg/ttl/ttlworker/job_manager_test.go b/pkg/ttl/ttlworker/job_manager_test.go index 240d1a4c94a38..68d1119eee054 100644 --- a/pkg/ttl/ttlworker/job_manager_test.go +++ b/pkg/ttl/ttlworker/job_manager_test.go @@ -16,6 +16,7 @@ package ttlworker import ( "context" + "encoding/json" "sync/atomic" "testing" "time" @@ -186,6 +187,90 @@ func newTTLTableStatusRows(status ...*cache.TableStatus) []chunk.Row { return rows } +func newTTLTaskRows(t *testing.T, tasks ...*cache.TTLTask) []chunk.Row { + c := chunk.NewChunkWithCapacity([]*types.FieldType{ + types.NewFieldType(mysql.TypeString), // job_id + types.NewFieldType(mysql.TypeLonglong), // table_id + types.NewFieldType(mysql.TypeLonglong), // scan_id + types.NewFieldType(mysql.TypeBlob), // scan_range_start + types.NewFieldType(mysql.TypeBlob), // scan_range_end + types.NewFieldType(mysql.TypeDatetime), // expire_time + types.NewFieldType(mysql.TypeString), // owner_id + types.NewFieldType(mysql.TypeString), // owner_addr + types.NewFieldType(mysql.TypeDatetime), // owner_hb_time + types.NewFieldType(mysql.TypeString), // status + types.NewFieldType(mysql.TypeDatetime), // status_update_time + types.NewFieldType(mysql.TypeString), // state + types.NewFieldType(mysql.TypeDatetime), // created_time + }, len(tasks)) + var rows []chunk.Row + + for _, task := range tasks { + jobID := types.NewDatum(task.JobID) + c.AppendDatum(0, &jobID) + tableID := types.NewDatum(task.TableID) + c.AppendDatum(1, &tableID) + scanID := types.NewDatum(task.ScanID) + c.AppendDatum(2, &scanID) + + if len(task.ScanRangeStart) == 0 { + c.AppendNull(3) + } else { + require.FailNow(t, "non-empty ScanRangeStart is not supported by this helper") + } + if len(task.ScanRangeEnd) == 0 { + c.AppendNull(4) + } else { + require.FailNow(t, "non-empty ScanRangeEnd is not supported by this helper") + } + + expireTime := types.NewDatum(types.NewTime(types.FromGoTime(task.ExpireTime), mysql.TypeDatetime, types.MaxFsp)) + c.AppendDatum(5, &expireTime) + + if task.OwnerID == "" { + c.AppendNull(6) + } else { + ownerID := types.NewDatum(task.OwnerID) + c.AppendDatum(6, &ownerID) + } + if task.OwnerAddr == "" { + c.AppendNull(7) + } else { + ownerAddr := types.NewDatum(task.OwnerAddr) + c.AppendDatum(7, &ownerAddr) + } + if task.OwnerHBTime.IsZero() { + c.AppendNull(8) + } else { + ownerHBTime := types.NewDatum(types.NewTime(types.FromGoTime(task.OwnerHBTime), mysql.TypeDatetime, types.MaxFsp)) + c.AppendDatum(8, &ownerHBTime) + } + + status := types.NewDatum(string(task.Status)) + c.AppendDatum(9, &status) + statusUpdateTime := types.NewDatum(types.NewTime(types.FromGoTime(task.StatusUpdateTime), mysql.TypeDatetime, types.MaxFsp)) + c.AppendDatum(10, &statusUpdateTime) + + if task.State == nil { + c.AppendNull(11) + } else { + stateJSON, err := json.Marshal(task.State) + require.NoError(t, err) + stateDatum := types.NewDatum(string(stateJSON)) + c.AppendDatum(11, &stateDatum) + } + + createdTime := types.NewDatum(types.NewTime(types.FromGoTime(task.CreatedTime), mysql.TypeDatetime, types.MaxFsp)) + c.AppendDatum(12, &createdTime) + } + + iter := chunk.NewIterator4Chunk(c) + for row := iter.Begin(); row != iter.End(); row = iter.Next() { + rows = append(rows, row) + } + return rows +} + var updateStatusSQL = "SELECT LOW_PRIORITY table_id,parent_table_id,table_statistics,last_job_id,last_job_start_time,last_job_finish_time,last_job_ttl_expire,last_job_summary,current_job_id,current_job_owner_id,current_job_owner_addr,current_job_owner_hb_time,current_job_start_time,current_job_ttl_expire,current_job_state,current_job_status,current_job_status_update_time FROM mysql.tidb_ttl_table_status" // TTLJob exports the ttlJob for test @@ -288,35 +373,100 @@ func (j *ttlJob) ID() string { } func TestCheckFinishedJobRecyclesExternalTTLTask(t *testing.T) { - createTime := time.Unix(1234, 0) - externalMgr := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} - m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) - m.runningJobs = []*ttlJob{ - { - id: "job1", - ownerID: "test-id", - createTime: createTime, - tableID: 1, - status: cache.JobStatusRunning, - }, - } + t.Run("all local jobs finish", func(t *testing.T) { + createTime := time.Unix(1234, 0) + externalMgr := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} + m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) + m.runningJobs = []*ttlJob{ + { + id: "job1", + ownerID: "test-id", + createTime: createTime, + tableID: 1, + status: cache.JobStatusRunning, + }, + } - se := newMockSession(t) - sqlCounter := 0 - se.executeSQL = func(_ context.Context, sql string, args ...any) ([]chunk.Row, error) { - sqlCounter++ - if sqlCounter == 1 { - expectedSQL, expectedArgs := cache.SelectFromTTLTaskWithJobID("job1") - require.Equal(t, expectedSQL, sql) - require.Equal(t, expectedArgs, args) + se := newMockSession(t) + sqlCounter := 0 + se.executeSQL = func(_ context.Context, sql string, args ...any) ([]chunk.Row, error) { + sqlCounter++ + if sqlCounter == 1 { + expectedSQL, expectedArgs := cache.SelectFromTTLTaskWithJobID("job1") + require.Equal(t, expectedSQL, sql) + require.Equal(t, expectedArgs, args) + } + return nil, nil } - return nil, nil - } - m.CheckFinishedJob(se) - require.Empty(t, m.runningJobs) - require.Equal(t, uint64(createTime.Unix()), externalMgr.recycledCreateTS) - require.Equal(t, 4, sqlCounter) + m.CheckFinishedJob(se) + require.Empty(t, m.runningJobs) + require.Equal(t, uint64(createTime.Unix()), externalMgr.recycledCreateTS) + require.Equal(t, 4, sqlCounter) + }) + + t.Run("recycle when only some local jobs finish", func(t *testing.T) { + finishedCreateTime := time.Unix(1234, 0) + runningCreateTime := time.Unix(2234, 0) + externalMgr := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} + m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) + m.runningJobs = []*ttlJob{ + { + id: "job-finished", + ownerID: "test-id", + createTime: finishedCreateTime, + tableID: 1, + status: cache.JobStatusRunning, + }, + { + id: "job-running", + ownerID: "test-id", + createTime: runningCreateTime, + tableID: 2, + status: cache.JobStatusRunning, + }, + } + + finishedTasks := newTTLTaskRows(t, &cache.TTLTask{ + JobID: "job-finished", + TableID: 1, + ScanID: 1, + ExpireTime: finishedCreateTime, + Status: cache.TaskStatusFinished, + StatusUpdateTime: finishedCreateTime, + CreatedTime: finishedCreateTime, + }) + runningTasks := newTTLTaskRows(t, &cache.TTLTask{ + JobID: "job-running", + TableID: 2, + ScanID: 1, + ExpireTime: runningCreateTime, + Status: cache.TaskStatusRunning, + StatusUpdateTime: runningCreateTime, + CreatedTime: runningCreateTime, + }) + + se := newMockSession(t) + se.executeSQL = func(_ context.Context, sql string, args ...any) ([]chunk.Row, error) { + expectedSQL, _ := cache.SelectFromTTLTaskWithJobID("job-finished") + if sql != expectedSQL { + return nil, nil + } + switch args[0] { + case "job-finished": + return finishedTasks, nil + case "job-running": + return runningTasks, nil + default: + return nil, nil + } + } + + m.CheckFinishedJob(se) + require.Len(t, m.runningJobs, 1) + require.Equal(t, "job-running", m.runningJobs[0].id) + require.Equal(t, uint64(finishedCreateTime.Unix()), externalMgr.recycledCreateTS) + }) } func TestCheckFinishedJobDoesNotRecycleExternalTTLTaskFromMaster(t *testing.T) { From 2f522693c4591e5c5d54314da6195e85b81bbcd4 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 15:05:36 +0800 Subject: [PATCH 12/16] domain: avoid local TTL fallback without controller --- pkg/domain/domain.go | 36 ++++++++++++++++++++++++++++-------- pkg/domain/domain_test.go | 13 +++++++++++++ 2 files changed, 41 insertions(+), 8 deletions(-) diff --git a/pkg/domain/domain.go b/pkg/domain/domain.go index 48cf1350d1ec0..379fd67ab1219 100644 --- a/pkg/domain/domain.go +++ b/pkg/domain/domain.go @@ -2892,9 +2892,13 @@ func (do *Domain) serverIDKeeper() { // StartTTLJobManager creates and starts the ttl job manager func (do *Domain) StartTTLJobManager() { + role, configured := do.ttlExternalWorkloadRole() if !do.shouldStartTTLJobManager() { - logutil.BgLogger().Info("don't run ttl job manager", - zap.String("role", string(do.extWorkloadMgr.Role()))) + fields := make([]zap.Field, 0, 1) + if configured { + fields = append(fields, zap.String("role", string(role))) + } + logutil.BgLogger().Info("don't run ttl job manager", fields...) return } ttlJobManager := ttlworker.NewJobManager(do.ddl.GetID(), do.advancedSysSessionPool, do.store, do.etcdClient, do.ddl.OwnerManager().IsOwner, do.extWorkloadMgr) @@ -2902,14 +2906,30 @@ func (do *Domain) StartTTLJobManager() { ttlJobManager.Start() } +func (do *Domain) ttlExternalWorkloadRole() (config.ExternalWorkloadRole, bool) { + if do.extWorkloadMgr != nil { + return do.extWorkloadMgr.Role(), true + } + cfg := config.GetGlobalConfig().ExternalWorkload + if !cfg.Enable { + return "", false + } + if cfg.Role == "" { + return config.RoleMaster, true + } + return cfg.Role, true +} + func (do *Domain) shouldStartTTLJobManager() bool { - // Unit tests and non-worker deployments do not install an external workload - // manager, so they still start a local TTL job manager. Once external - // workload is enabled, only the dedicated TTL task worker should run TTL jobs. - if extworkload.IsEnabled(do.extWorkloadMgr) && !extworkload.IsTTLTaskWorker(do.extWorkloadMgr) { - return false + // Once external workload is configured, TTL jobs must run only on the + // dedicated TTL task worker with a live controller manager. Falling back to + // local TTL scheduling when controller coordination is unavailable can cause + // split-brain ownership with externally assigned TTL tasks. + role, configured := do.ttlExternalWorkloadRole() + if !configured { + return true } - return true + return do.extWorkloadMgr != nil && role == config.RoleTTLTaskWorker } // TTLJobManager returns the ttl job manager on this domain diff --git a/pkg/domain/domain_test.go b/pkg/domain/domain_test.go index 68141377ff174..a31c84190137c 100644 --- a/pkg/domain/domain_test.go +++ b/pkg/domain/domain_test.go @@ -280,6 +280,7 @@ func TestUpdateExternalWorkloadTTLJobEnableOnlyFromMaster(t *testing.T) { } func TestShouldStartTTLJobManagerWithExternalWorkloadRole(t *testing.T) { + t.Cleanup(config.RestoreFunc()) dom := NewMockDomain() require.True(t, dom.shouldStartTTLJobManager()) @@ -288,6 +289,18 @@ func TestShouldStartTTLJobManagerWithExternalWorkloadRole(t *testing.T) { dom.SetExternalWorkloadManager(&fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker}) require.True(t, dom.shouldStartTTLJobManager()) + + dom.SetExternalWorkloadManager(nil) + config.UpdateGlobal(func(conf *config.Config) { + conf.ExternalWorkload.Enable = true + conf.ExternalWorkload.Role = config.RoleMaster + }) + require.False(t, dom.shouldStartTTLJobManager()) + + config.UpdateGlobal(func(conf *config.Config) { + conf.ExternalWorkload.Role = config.RoleTTLTaskWorker + }) + require.False(t, dom.shouldStartTTLJobManager()) } // ETCD use ip:port as unix socket address, however this address is invalid on windows. From 772fd519447aa1c6f9f1a4a455429d0e4d7a6cb7 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Wed, 29 Jul 2026 15:10:03 +0800 Subject: [PATCH 13/16] ddl: register TTL tables in batch create jobs --- pkg/ddl/create_table.go | 3 +++ pkg/ddl/db_table_test.go | 25 +++++++++++++++++++++++++ 2 files changed, 28 insertions(+) diff --git a/pkg/ddl/create_table.go b/pkg/ddl/create_table.go index 7c302a7dc1e4b..38a85e7409d5f 100644 --- a/pkg/ddl/create_table.go +++ b/pkg/ddl/create_table.go @@ -355,6 +355,9 @@ func (w *worker) onCreateTables(jobCtx *jobContext, job *model.Job) (int64, erro return ver, errors.Trace(err) } } + for i := range tableInfos { + w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tableInfos[i]) + } job.State = model.JobStateDone job.SchemaState = model.StatePublic diff --git a/pkg/ddl/db_table_test.go b/pkg/ddl/db_table_test.go index 8a1c3f3a8b00e..7114d125f1180 100644 --- a/pkg/ddl/db_table_test.go +++ b/pkg/ddl/db_table_test.go @@ -414,6 +414,31 @@ func TestBatchCreateTable(t *testing.T) { tk.Session().SetValue(sessionctx.QueryString, "skip") err = d.BatchCreateTableWithInfo(tk.Session(), ast.NewCIStr("test"), []*model.TableInfo{newinfo}, ddl.WithOnExist(ddl.OnExistError)) require.NoError(t, err) + + t.Run("batch create ttl tables registers external workload", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, store := createTTLExternalWorkloadTestKit(t, mgr) + d := domain.GetDomain(tk.Session()).DDLExecutor() + infos := make([]*model.TableInfo, 0, 2) + for _, name := range []string{"ttl_batch_1", "ttl_batch_2"} { + info, err := testTableInfo(store, name, 2) + require.NoError(t, err) + info.Columns[0].FieldType = *types.NewFieldType(mysql.TypeDatetime) + info.TTLInfo = &model.TTLInfo{ + ColumnName: info.Columns[0].Name, + IntervalExprStr: "1", + IntervalTimeUnit: int(ast.TimeUnitDay), + Enable: true, + JobInterval: model.DefaultTTLJobInterval, + } + infos = append(infos, info) + } + + tk.Session().SetValue(sessionctx.QueryString, "skip") + err := d.BatchCreateTableWithInfo(tk.Session(), ast.NewCIStr("test"), infos, ddl.WithOnExist(ddl.OnExistError), ddl.WithIDAllocated(true)) + require.NoError(t, err) + require.Equal(t, []int64{infos[0].ID, infos[1].ID}, mgr.registeredTTLTables()) + }) } // port from mysql From 343ca684dbac62ee466b824eebfb1346dacd4aaa Mon Sep 17 00:00:00 2001 From: ystaticy Date: Fri, 31 Jul 2026 13:46:44 +0800 Subject: [PATCH 14/16] fix comment Signed-off-by: ystaticy --- pkg/ttl/ttlworker/job_manager.go | 3 +- pkg/ttl/ttlworker/job_manager_test.go | 60 +++++++++++++-------------- 2 files changed, 32 insertions(+), 31 deletions(-) diff --git a/pkg/ttl/ttlworker/job_manager.go b/pkg/ttl/ttlworker/job_manager.go index e2386c8cf9f4a..23dbedb64eb54 100644 --- a/pkg/ttl/ttlworker/job_manager.go +++ b/pkg/ttl/ttlworker/job_manager.go @@ -602,6 +602,7 @@ func (m *JobManager) findAllTasksForJob(se session.Session, jobID string) ([]*ca } func (m *JobManager) checkFinishedJob(se session.Session) { + runningJobsCount := len(m.runningJobs) totalFinishedJobs := 0 maxJobCreateTime := uint64(0) // reverse iteration so that we could remove the job safely in the loop @@ -643,7 +644,7 @@ func (m *JobManager) checkFinishedJob(se session.Session) { } } } - if totalFinishedJobs > 0 && extworkload.IsTTLTaskWorker(m.extWorkload) { + if runningJobsCount > 0 && totalFinishedJobs == runningJobsCount && extworkload.IsTTLTaskWorker(m.extWorkload) { if err := m.extWorkload.RecycleTTLTask(m.ctx, maxJobCreateTime); err != nil { logutil.Logger(m.ctx).Warn("failed to recycle TTL task from external workload controller", zap.Uint64("completedJobCreateTime", maxJobCreateTime), diff --git a/pkg/ttl/ttlworker/job_manager_test.go b/pkg/ttl/ttlworker/job_manager_test.go index 68d1119eee054..4427adfa17fbe 100644 --- a/pkg/ttl/ttlworker/job_manager_test.go +++ b/pkg/ttl/ttlworker/job_manager_test.go @@ -405,46 +405,46 @@ func TestCheckFinishedJobRecyclesExternalTTLTask(t *testing.T) { require.Equal(t, 4, sqlCounter) }) - t.Run("recycle when only some local jobs finish", func(t *testing.T) { - finishedCreateTime := time.Unix(1234, 0) - runningCreateTime := time.Unix(2234, 0) + t.Run("do not recycle when an older local job is still running", func(t *testing.T) { + runningCreateTime := time.Unix(1234, 0) + finishedCreateTime := time.Unix(2234, 0) externalMgr := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} m := NewJobManager("test-id", nil, nil, nil, nil, externalMgr) m.runningJobs = []*ttlJob{ { - id: "job-finished", + id: "job-running", ownerID: "test-id", - createTime: finishedCreateTime, + createTime: runningCreateTime, tableID: 1, status: cache.JobStatusRunning, }, { - id: "job-running", + id: "job-finished", ownerID: "test-id", - createTime: runningCreateTime, + createTime: finishedCreateTime, tableID: 2, status: cache.JobStatusRunning, }, } - finishedTasks := newTTLTaskRows(t, &cache.TTLTask{ - JobID: "job-finished", - TableID: 1, - ScanID: 1, - ExpireTime: finishedCreateTime, - Status: cache.TaskStatusFinished, - StatusUpdateTime: finishedCreateTime, - CreatedTime: finishedCreateTime, - }) runningTasks := newTTLTaskRows(t, &cache.TTLTask{ JobID: "job-running", - TableID: 2, + TableID: 1, ScanID: 1, ExpireTime: runningCreateTime, Status: cache.TaskStatusRunning, StatusUpdateTime: runningCreateTime, CreatedTime: runningCreateTime, }) + finishedTasks := newTTLTaskRows(t, &cache.TTLTask{ + JobID: "job-finished", + TableID: 2, + ScanID: 1, + ExpireTime: finishedCreateTime, + Status: cache.TaskStatusFinished, + StatusUpdateTime: finishedCreateTime, + CreatedTime: finishedCreateTime, + }) se := newMockSession(t) se.executeSQL = func(_ context.Context, sql string, args ...any) ([]chunk.Row, error) { @@ -453,21 +453,21 @@ func TestCheckFinishedJobRecyclesExternalTTLTask(t *testing.T) { return nil, nil } switch args[0] { - case "job-finished": - return finishedTasks, nil - case "job-running": - return runningTasks, nil - default: - return nil, nil + case "job-running": + return runningTasks, nil + case "job-finished": + return finishedTasks, nil + default: + return nil, nil + } } - } - m.CheckFinishedJob(se) - require.Len(t, m.runningJobs, 1) - require.Equal(t, "job-running", m.runningJobs[0].id) - require.Equal(t, uint64(finishedCreateTime.Unix()), externalMgr.recycledCreateTS) - }) -} + m.CheckFinishedJob(se) + require.Len(t, m.runningJobs, 1) + require.Equal(t, "job-running", m.runningJobs[0].id) + require.Equal(t, uint64(0), externalMgr.recycledCreateTS) + }) + } func TestCheckFinishedJobDoesNotRecycleExternalTTLTaskFromMaster(t *testing.T) { externalMgr := &fakeExternalWorkloadManager{role: config.RoleMaster} From 98f5e3bd6e3b8237e7897b69c1067b760116f9c6 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Fri, 31 Jul 2026 14:58:41 +0800 Subject: [PATCH 15/16] fix ci Signed-off-by: ystaticy --- pkg/ttl/ttlworker/job_manager_test.go | 26 +++++++++++++------------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/pkg/ttl/ttlworker/job_manager_test.go b/pkg/ttl/ttlworker/job_manager_test.go index 4427adfa17fbe..3d535dfeb9ae1 100644 --- a/pkg/ttl/ttlworker/job_manager_test.go +++ b/pkg/ttl/ttlworker/job_manager_test.go @@ -453,21 +453,21 @@ func TestCheckFinishedJobRecyclesExternalTTLTask(t *testing.T) { return nil, nil } switch args[0] { - case "job-running": - return runningTasks, nil - case "job-finished": - return finishedTasks, nil - default: - return nil, nil - } + case "job-running": + return runningTasks, nil + case "job-finished": + return finishedTasks, nil + default: + return nil, nil } + } - m.CheckFinishedJob(se) - require.Len(t, m.runningJobs, 1) - require.Equal(t, "job-running", m.runningJobs[0].id) - require.Equal(t, uint64(0), externalMgr.recycledCreateTS) - }) - } + m.CheckFinishedJob(se) + require.Len(t, m.runningJobs, 1) + require.Equal(t, "job-running", m.runningJobs[0].id) + require.Equal(t, uint64(0), externalMgr.recycledCreateTS) + }) +} func TestCheckFinishedJobDoesNotRecycleExternalTTLTaskFromMaster(t *testing.T) { externalMgr := &fakeExternalWorkloadManager{role: config.RoleMaster} From 60afe2b101e68eb342b9fd26db53078b782d1322 Mon Sep 17 00:00:00 2001 From: ystaticy Date: Sun, 2 Aug 2026 14:45:57 +0800 Subject: [PATCH 16/16] ttl: fail DDL on external TTL sync errors Signed-off-by: ystaticy --- pkg/ddl/BUILD.bazel | 2 + pkg/ddl/create_table.go | 24 +++++- pkg/ddl/db_table_test.go | 34 ++++++++ pkg/ddl/external_workload_ttl_test.go | 115 ++++++++++++++++++++++++++ pkg/ddl/table.go | 40 ++++++--- pkg/ddl/ttl.go | 45 +++++----- pkg/ddl/ttl_test.go | 24 +++--- 7 files changed, 234 insertions(+), 50 deletions(-) diff --git a/pkg/ddl/BUILD.bazel b/pkg/ddl/BUILD.bazel index 25cb25d32e669..73c29c89a2a7e 100644 --- a/pkg/ddl/BUILD.bazel +++ b/pkg/ddl/BUILD.bazel @@ -419,6 +419,7 @@ go_test( "@com_github_pingcap_errors//:errors", "@com_github_pingcap_failpoint//:failpoint", "@com_github_pingcap_kvproto//pkg/keyspacepb", + "@com_github_pingcap_log//:log", "@com_github_prometheus_client_golang//prometheus", "@com_github_prometheus_client_model//go", "@com_github_stretchr_testify//assert", @@ -436,5 +437,6 @@ go_test( "@org_uber_go_goleak//:goleak", "@org_uber_go_mock//gomock", "@org_uber_go_zap//:zap", + "@org_uber_go_zap//zaptest/observer", ], ) diff --git a/pkg/ddl/create_table.go b/pkg/ddl/create_table.go index 38a85e7409d5f..1dd905b1612a5 100644 --- a/pkg/ddl/create_table.go +++ b/pkg/ddl/create_table.go @@ -249,7 +249,9 @@ func (w *worker) onCreateTable(jobCtx *jobContext, job *model.Job) (ver int64, _ return ver, errors.Trace(err) } - w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo) + if err := w.registerTTLTableToExternalWorkload(jobCtx.ctx, tbInfo); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } // Finish this job. job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo) @@ -295,7 +297,9 @@ func (w *worker) createTableWithForeignKeys(jobCtx *jobContext, job *model.Job, if err != nil { return ver, errors.Trace(err) } - w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo) + if err := w.registerTTLTableToExternalWorkload(jobCtx.ctx, tbInfo); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo) return ver, nil @@ -355,8 +359,22 @@ func (w *worker) onCreateTables(jobCtx *jobContext, job *model.Job) (int64, erro return ver, errors.Trace(err) } } + registeredTTLTableIDs := make([]int64, 0, len(tableInfos)) for i := range tableInfos { - w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tableInfos[i]) + if err := w.registerTTLTableToExternalWorkload(jobCtx.ctx, tableInfos[i]); err != nil { + for j := len(registeredTTLTableIDs) - 1; j >= 0; j-- { + compensateTableID := registeredTTLTableIDs[j] + if compensateErr := w.deleteTTLTableFromExternalWorkload(jobCtx.ctx, compensateTableID); compensateErr != nil { + logutil.DDLLogger().Warn("failed to roll back TTL table registration in external workload controller", + zap.Int64("tableID", compensateTableID), + zap.Error(compensateErr)) + } + } + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } + if tableInfos[i] != nil && tableInfos[i].TTLInfo != nil && tableInfos[i].TTLInfo.Enable { + registeredTTLTableIDs = append(registeredTTLTableIDs, tableInfos[i].ID) + } } job.State = model.JobStateDone diff --git a/pkg/ddl/db_table_test.go b/pkg/ddl/db_table_test.go index 7114d125f1180..8ce1acc78483b 100644 --- a/pkg/ddl/db_table_test.go +++ b/pkg/ddl/db_table_test.go @@ -439,6 +439,40 @@ func TestBatchCreateTable(t *testing.T) { require.NoError(t, err) require.Equal(t, []int64{infos[0].ID, infos[1].ID}, mgr.registeredTTLTables()) }) + + t.Run("batch create ttl tables unregisters previous external registrations on failure", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, store := createTTLExternalWorkloadTestKit(t, mgr) + d := domain.GetDomain(tk.Session()).DDLExecutor() + infos := make([]*model.TableInfo, 0, 2) + for _, name := range []string{"ttl_batch_fail_1", "ttl_batch_fail_2"} { + info, err := testTableInfo(store, name, 2) + require.NoError(t, err) + info.Columns[0].FieldType = *types.NewFieldType(mysql.TypeDatetime) + info.TTLInfo = &model.TTLInfo{ + ColumnName: info.Columns[0].Name, + IntervalExprStr: "1", + IntervalTimeUnit: int(ast.TimeUnitDay), + Enable: true, + JobInterval: model.DefaultTTLJobInterval, + } + infos = append(infos, info) + } + mgr.registerErrFn = func(tableID int64) error { + if tableID == infos[1].ID { + return context.DeadlineExceeded + } + return nil + } + + tk.Session().SetValue(sessionctx.QueryString, "skip") + err := d.BatchCreateTableWithInfo(tk.Session(), ast.NewCIStr("test"), infos, ddl.WithOnExist(ddl.OnExistError), ddl.WithIDAllocated(true)) + require.ErrorContains(t, err, context.DeadlineExceeded.Error()) + require.Equal(t, []int64{infos[0].ID}, mgr.registeredTTLTables()) + require.Equal(t, []int64{infos[0].ID}, mgr.deletedTTLTables()) + tk.MustQuery("show tables like 'ttl_batch_fail_1'").Check(testkit.Rows()) + tk.MustQuery("show tables like 'ttl_batch_fail_2'").Check(testkit.Rows()) + }) } // port from mysql diff --git a/pkg/ddl/external_workload_ttl_test.go b/pkg/ddl/external_workload_ttl_test.go index 493b258a03503..6b34298aa2be6 100644 --- a/pkg/ddl/external_workload_ttl_test.go +++ b/pkg/ddl/external_workload_ttl_test.go @@ -21,6 +21,7 @@ import ( "time" "github.com/pingcap/kvproto/pkg/keyspacepb" + "github.com/pingcap/log" "github.com/pingcap/tidb/pkg/config" "github.com/pingcap/tidb/pkg/domain" "github.com/pingcap/tidb/pkg/kv" @@ -30,6 +31,8 @@ import ( "github.com/pingcap/tidb/pkg/testkit" "github.com/pingcap/tidb/pkg/testkit/external" "github.com/stretchr/testify/require" + "go.uber.org/zap" + "go.uber.org/zap/zaptest/observer" ) type recordingExternalWorkloadManager struct { @@ -38,6 +41,8 @@ type recordingExternalWorkloadManager struct { mu sync.Mutex registeredTables []int64 deletedTables []int64 + registerErrFn func(int64) error + deleteErrFn func(int64) error } func (*recordingExternalWorkloadManager) Close() error { return nil } @@ -69,6 +74,11 @@ func (*recordingExternalWorkloadManager) UpdateGCLifeTime(context.Context, time. func (m *recordingExternalWorkloadManager) RegisterTTLTask(_ context.Context, tableID int64, _ bool) error { m.mu.Lock() defer m.mu.Unlock() + if m.registerErrFn != nil { + if err := m.registerErrFn(tableID); err != nil { + return err + } + } m.registeredTables = append(m.registeredTables, tableID) return nil } @@ -76,6 +86,11 @@ func (m *recordingExternalWorkloadManager) RegisterTTLTask(_ context.Context, ta func (m *recordingExternalWorkloadManager) DeleteTTLTableInfo(_ context.Context, tableID int64) error { m.mu.Lock() defer m.mu.Unlock() + if m.deleteErrFn != nil { + if err := m.deleteErrFn(tableID); err != nil { + return err + } + } m.deletedTables = append(m.deletedTables, tableID) return nil } @@ -134,6 +149,22 @@ func createTTLExternalWorkloadTestKit(t *testing.T, mgr *recordingExternalWorklo } func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { + t.Run("create table registration failure aborts ddl", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{ + role: config.RoleMaster, + registerErrFn: func(int64) error { return context.DeadlineExceeded }, + } + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + err := tk.ExecToErr(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + require.ErrorContains(t, err, context.DeadlineExceeded.Error()) + tk.MustQuery("show tables like 't'").Check(testkit.Rows()) + require.Empty(t, mgr.registeredTTLTables()) + }) + t.Run("create table with foreign key registers ttl", func(t *testing.T) { mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} tk, _ := createTTLExternalWorkloadTestKit(t, mgr) @@ -167,6 +198,31 @@ func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { require.Equal(t, []int64{tbl.Meta().ID}, mgr.deletedTTLTables()) }) + t.Run("drop table delete failure aborts ddl", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + tbl := external.GetTableByName(t, tk, "test", "t") + mgr.deleteErrFn = func(tableID int64) error { + if tableID == tbl.Meta().ID { + return context.DeadlineExceeded + } + return nil + } + + err := tk.ExecToErr("drop table t") + require.ErrorContains(t, err, context.DeadlineExceeded.Error()) + + currentTbl := external.GetTableByName(t, tk, "test", "t") + require.Equal(t, tbl.Meta().ID, currentTbl.Meta().ID) + require.Empty(t, mgr.deletedTTLTables()) + }) + t.Run("truncate table refreshes ttl metadata", func(t *testing.T) { mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} tk, _ := createTTLExternalWorkloadTestKit(t, mgr) @@ -185,6 +241,65 @@ func TestExternalWorkloadTTLDDLIntegration(t *testing.T) { require.Equal(t, []int64{oldTbl.Meta().ID, newTbl.Meta().ID}, mgr.registeredTTLTables()) }) + t.Run("truncate table register failure restores old ttl registration", func(t *testing.T) { + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + oldTbl := external.GetTableByName(t, tk, "test", "t") + mgr.registerErrFn = func(tableID int64) error { + if tableID != oldTbl.Meta().ID { + return context.DeadlineExceeded + } + return nil + } + + err := tk.ExecToErr("truncate table t") + require.ErrorContains(t, err, context.DeadlineExceeded.Error()) + + currentTbl := external.GetTableByName(t, tk, "test", "t") + require.Equal(t, oldTbl.Meta().ID, currentTbl.Meta().ID) + require.Equal(t, []int64{oldTbl.Meta().ID}, mgr.deletedTTLTables()) + require.Equal(t, []int64{oldTbl.Meta().ID, oldTbl.Meta().ID}, mgr.registeredTTLTables()) + }) + + t.Run("truncate table compensation failure logs keyword", func(t *testing.T) { + core, recorded := observer.New(zap.WarnLevel) + restoreLog := log.ReplaceGlobals(zap.New(core), &log.ZapProperties{Level: zap.NewAtomicLevelAt(zap.InfoLevel)}) + defer restoreLog() + + mgr := &recordingExternalWorkloadManager{role: config.RoleMaster} + tk, _ := createTTLExternalWorkloadTestKit(t, mgr) + + tk.MustExec(`create table t( + id int primary key, + created_at datetime + ) TTL = created_at + interval 1 day`) + + oldTbl := external.GetTableByName(t, tk, "test", "t") + mgr.registerErrFn = func(int64) error { return context.DeadlineExceeded } + + err := tk.ExecToErr("truncate table t") + require.ErrorContains(t, err, context.DeadlineExceeded.Error()) + + found := false + for _, entry := range recorded.All() { + if entry.Message != "truncate TTL external workload compensation failed" { + continue + } + fields := entry.ContextMap() + require.Equal(t, "truncate_ttl_restore_old_registration_failed", fields["keyword"]) + require.EqualValues(t, oldTbl.Meta().ID, fields["oldTableID"]) + found = true + break + } + require.True(t, found) + }) + t.Run("ddl syncs ttl metadata from ttl worker role", func(t *testing.T) { mgr := &recordingExternalWorkloadManager{role: config.RoleTTLTaskWorker} tk, _ := createTTLExternalWorkloadTestKit(t, mgr) diff --git a/pkg/ddl/table.go b/pkg/ddl/table.go index 1e0da1b67c5ab..2f6ebc37cd311 100644 --- a/pkg/ddl/table.go +++ b/pkg/ddl/table.go @@ -50,6 +50,8 @@ import ( const tiflashCheckTiDBHTTPAPIHalfInterval = 2500 * time.Millisecond +const truncateTTLOldRegistrationCompensationLogKeyword = "truncate_ttl_restore_old_registration_failed" + func repairTableOrViewWithCheck(t *meta.Mutator, job *model.Job, schemaID int64, tbInfo *model.TableInfo) error { err := checkTableInfoValid(tbInfo) if err != nil { @@ -81,6 +83,11 @@ func (w *worker) onDropTableOrView(jobCtx *jobContext, job *model.Job) (ver int6 return ver, err } } + if jobCtx.oldDDLCtx != nil && tblInfo.TTLInfo != nil { + if err := jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } + } tblInfo.State = model.StateWriteOnly ver, err = updateVersionAndTableInfo(jobCtx, job, tblInfo, originalState != tblInfo.State) if err != nil { @@ -140,9 +147,6 @@ func (w *worker) onDropTableOrView(jobCtx *jobContext, job *model.Job) (ver int6 if err := w.dropMaskingPoliciesOnTable(jobCtx, tblInfo.ID); err != nil { return ver, errors.Wrapf(err, "failed to drop masking policies on table %d", tblInfo.ID) } - if jobCtx.oldDDLCtx != nil && tblInfo.TTLInfo != nil { - jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID) - } // Finish this job. job.FinishTableJob(model.JobStateDone, model.StateNone, ver, tblInfo) @@ -478,6 +482,28 @@ func (w *worker) onTruncateTable(jobCtx *jobContext, job *model.Job) (ver int64, if err != nil { return ver, err } + if jobCtx.oldDDLCtx != nil { + if oldTblInfo.TTLInfo != nil { + if err := jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, oldTblInfo.ID); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } + } + if oldTblInfo.TTLInfo != nil && oldTblInfo.TTLInfo.Enable { + newTTLTableInfo := oldTblInfo.Clone() + newTTLTableInfo.ID = args.NewTableID + if err := jobCtx.oldDDLCtx.registerTTLTableToExternalWorkload(jobCtx.ctx, newTTLTableInfo); err != nil { + if compensateErr := jobCtx.oldDDLCtx.registerTTLTableToExternalWorkload(jobCtx.ctx, oldTblInfo); compensateErr != nil { + logutil.DDLLogger().Warn("truncate TTL external workload compensation failed", + zap.String("keyword", truncateTTLOldRegistrationCompensationLogKeyword), + zap.Int64("oldTableID", oldTblInfo.ID), + zap.Int64("newTableID", newTTLTableInfo.ID), + zap.Error(compensateErr), + zap.NamedError("registerNewTableErr", err)) + } + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } + } + } err = metaMut.DropTableOrView(schemaID, tblInfo.ID) if err != nil { job.State = model.JobStateCancelled @@ -629,14 +655,6 @@ func (w *worker) onTruncateTable(jobCtx *jobContext, job *model.Job) (ver int64, if err != nil { return ver, errors.Trace(err) } - if jobCtx.oldDDLCtx != nil { - if oldTblInfo.TTLInfo != nil { - jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, oldTblInfo.ID) - } - if tblInfo.TTLInfo != nil && tblInfo.TTLInfo.Enable { - jobCtx.oldDDLCtx.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tblInfo) - } - } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) // see truncateTableByReassignPartitionIDs for why they might change. diff --git a/pkg/ddl/ttl.go b/pkg/ddl/ttl.go index a5169bb010ab5..5fcba4d726d1b 100644 --- a/pkg/ddl/ttl.go +++ b/pkg/ddl/ttl.go @@ -20,7 +20,6 @@ import ( "time" "github.com/pingcap/errors" - "github.com/pingcap/tidb/pkg/ddl/logutil" "github.com/pingcap/tidb/pkg/extworkload" infoschemactx "github.com/pingcap/tidb/pkg/infoschema/context" "github.com/pingcap/tidb/pkg/meta/model" @@ -31,7 +30,6 @@ import ( "github.com/pingcap/tidb/pkg/ttl/cache" "github.com/pingcap/tidb/pkg/types" "github.com/pingcap/tidb/pkg/util/dbterror" - "go.uber.org/zap" ) func onTTLInfoRemove(jobCtx *jobContext, job *model.Job) (ver int64, err error) { @@ -46,7 +44,9 @@ func onTTLInfoRemove(jobCtx *jobContext, job *model.Job) (ver int64, err error) return ver, errors.Trace(err) } if jobCtx.oldDDLCtx != nil { - jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID) + if err := jobCtx.oldDDLCtx.deleteTTLTableFromExternalWorkload(jobCtx.ctx, tblInfo.ID); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) return ver, nil @@ -96,7 +96,9 @@ func onTTLInfoChange(jobCtx *jobContext, job *model.Job) (ver int64, err error) return ver, errors.Trace(err) } if jobCtx.oldDDLCtx != nil { - jobCtx.oldDDLCtx.syncTTLTableToExternalWorkload(jobCtx.ctx, tblInfo) + if err := jobCtx.oldDDLCtx.syncTTLTableToExternalWorkload(jobCtx.ctx, tblInfo); err != nil { + return ver, cancelJobOnExternalTTLWorkloadError(job, err) + } } job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tblInfo) return ver, nil @@ -110,6 +112,14 @@ func (dc *ddlCtx) externalWorkloadManager() (extworkload.Manager, bool) { return manager, extworkload.IsEnabled(manager) } +func cancelJobOnExternalTTLWorkloadError(job *model.Job, err error) error { + if err == nil { + return nil + } + job.State = model.JobStateCancelled + return errors.Trace(err) +} + func (dc *ddlCtx) registerTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) error { manager, ok := dc.externalWorkloadManager() if !ok || tblInfo == nil || tblInfo.TTLInfo == nil || !tblInfo.TTLInfo.Enable { @@ -118,35 +128,22 @@ func (dc *ddlCtx) registerTTLTableToExternalWorkload(ctx context.Context, tblInf return manager.RegisterTTLTask(ctx, tblInfo.ID, vardef.EnableTTLJob.Load()) } -func (dc *ddlCtx) tryRegisterTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) { - if err := dc.registerTTLTableToExternalWorkload(ctx, tblInfo); err != nil { - logutil.DDLLogger().Warn("failed to register TTL table to external workload controller", - zap.Int64("tableID", tblInfo.ID), - zap.Error(err)) - } -} - -func (dc *ddlCtx) syncTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) { +func (dc *ddlCtx) syncTTLTableToExternalWorkload(ctx context.Context, tblInfo *model.TableInfo) error { if tblInfo == nil { - return + return nil } if tblInfo.TTLInfo == nil || !tblInfo.TTLInfo.Enable { - dc.deleteTTLTableFromExternalWorkload(ctx, tblInfo.ID) - return + return dc.deleteTTLTableFromExternalWorkload(ctx, tblInfo.ID) } - dc.tryRegisterTTLTableToExternalWorkload(ctx, tblInfo) + return dc.registerTTLTableToExternalWorkload(ctx, tblInfo) } -func (dc *ddlCtx) deleteTTLTableFromExternalWorkload(ctx context.Context, tableID int64) { +func (dc *ddlCtx) deleteTTLTableFromExternalWorkload(ctx context.Context, tableID int64) error { manager, ok := dc.externalWorkloadManager() if !ok { - return - } - if err := manager.DeleteTTLTableInfo(ctx, tableID); err != nil { - logutil.DDLLogger().Warn("failed to delete TTL table from external workload controller", - zap.Int64("tableID", tableID), - zap.Error(err)) + return nil } + return manager.DeleteTTLTableInfo(ctx, tableID) } // checkTTLInfoValid checks the TTL settings for a table. diff --git a/pkg/ddl/ttl_test.go b/pkg/ddl/ttl_test.go index 712acbac57f15..504b17026f061 100644 --- a/pkg/ddl/ttl_test.go +++ b/pkg/ddl/ttl_test.go @@ -161,6 +161,7 @@ type fakeExternalWorkloadManager struct { registeredTable int64 registerEnabled bool deletedTable int64 + deleteErr error recycledCreateTS uint64 updatedEnable *bool registerErr error @@ -191,7 +192,7 @@ func (m *fakeExternalWorkloadManager) RegisterTTLTask(_ context.Context, tableID } func (m *fakeExternalWorkloadManager) DeleteTTLTableInfo(_ context.Context, tableID int64) error { m.deletedTable = tableID - return nil + return m.deleteErr } func (m *fakeExternalWorkloadManager) RecycleTTLTask(_ context.Context, completedJobCreateTime uint64) error { m.recycledCreateTS = completedJobCreateTime @@ -226,13 +227,13 @@ func TestExternalWorkloadTTLTableReportsOnlyFromMaster(t *testing.T) { require.Equal(t, int64(123), master.registeredTable) require.False(t, master.registerEnabled) - dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID) + require.NoError(t, dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID)) require.Equal(t, int64(123), master.deletedTable) ttlWorker := &fakeExternalWorkloadManager{role: config.RoleTTLTaskWorker} dc = &ddlCtx{extWorkload: ttlWorker} require.NoError(t, dc.registerTTLTableToExternalWorkload(context.Background(), tblInfo)) - dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID) + require.NoError(t, dc.deleteTTLTableFromExternalWorkload(context.Background(), tblInfo.ID)) require.Equal(t, int64(123), ttlWorker.registeredTable) require.Equal(t, int64(123), ttlWorker.deletedTable) } @@ -258,22 +259,21 @@ func TestExternalWorkloadTTLTableRegisterReturnsError(t *testing.T) { require.ErrorIs(t, err, boom) } -func TestExternalWorkloadTTLTableTryRegisterSwallowsError(t *testing.T) { - manager := &fakeExternalWorkloadManager{role: config.RoleMaster, registerErr: errors.New("boom")} +func TestExternalWorkloadTTLTableDeleteReturnsError(t *testing.T) { + boom := errors.New("boom") + manager := &fakeExternalWorkloadManager{role: config.RoleMaster, deleteErr: boom} dc := &ddlCtx{extWorkload: manager} - dc.tryRegisterTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ - ID: 123, - TTLInfo: &model.TTLInfo{Enable: true}, - }) - require.Equal(t, int64(123), manager.registeredTable) + err := dc.deleteTTLTableFromExternalWorkload(context.Background(), 123) + require.ErrorIs(t, err, boom) + require.Equal(t, int64(123), manager.deletedTable) } func TestExternalWorkloadTTLTableSyncDeletesDisabledTTL(t *testing.T) { manager := &fakeExternalWorkloadManager{role: config.RoleMaster} dc := &ddlCtx{extWorkload: manager} - dc.syncTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ + require.NoError(t, dc.syncTTLTableToExternalWorkload(context.Background(), &model.TableInfo{ ID: 123, TTLInfo: &model.TTLInfo{Enable: false}, - }) + })) require.Equal(t, int64(123), manager.deletedTable) }