diff --git a/pkg/utils/etcdutil/etcdutil.go b/pkg/utils/etcdutil/etcdutil.go index 34a54738f5..5592cc5cc2 100644 --- a/pkg/utils/etcdutil/etcdutil.go +++ b/pkg/utils/etcdutil/etcdutil.go @@ -161,7 +161,19 @@ func RemoveEtcdMember(client *clientv3.Client, id uint64) (*clientv3.MemberRemov // EtcdKVGet returns the etcd GetResponse by given key or key prefix func EtcdKVGet(c *clientv3.Client, key string, opts ...clientv3.OpOption) (*clientv3.GetResponse, error) { - ctx, cancel := context.WithTimeout(c.Ctx(), DefaultRequestTimeout) + return EtcdKVGetWithContext(c.Ctx(), c, key, opts...) +} + +// EtcdKVGetWithContext returns the etcd GetResponse by the given key or key +// prefix. The request stops when either ctx is canceled or the default request +// timeout is reached. +func EtcdKVGetWithContext( + ctx context.Context, + c *clientv3.Client, + key string, + opts ...clientv3.OpOption, +) (*clientv3.GetResponse, error) { + ctx, cancel := context.WithTimeout(ctx, DefaultRequestTimeout) defer cancel() start := time.Now() diff --git a/server/microservice_cleanup.go b/server/microservice_cleanup.go new file mode 100644 index 0000000000..fb28f3a7ae --- /dev/null +++ b/server/microservice_cleanup.go @@ -0,0 +1,217 @@ +// Copyright 2026 TiKV Project Authors. +// +// 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 server + +import ( + "context" + "encoding/json" + goerrors "errors" + "time" + + clientv3 "go.etcd.io/etcd/client/v3" + "go.uber.org/zap" + + "github.com/pingcap/errors" + "github.com/pingcap/failpoint" + "github.com/pingcap/log" + + "github.com/tikv/pd/pkg/errs" + "github.com/tikv/pd/pkg/keyspace/constant" + "github.com/tikv/pd/pkg/storage/endpoint" + "github.com/tikv/pd/pkg/storage/kv" + "github.com/tikv/pd/pkg/utils/etcdutil" + "github.com/tikv/pd/pkg/utils/keypath" + "github.com/tikv/pd/pkg/utils/logutil" +) + +const ( + microserviceMetadataCleanupInitialRetryInterval = time.Second + microserviceMetadataCleanupMaxRetryInterval = 30 * time.Second +) + +var errMicroserviceMetadataCleanupRejected = errors.New("microservice metadata cleanup rejected") + +type microserviceMetadataCleanupTerm struct { + leaderKey string + leaderValue string + leaseID clientv3.LeaseID +} + +func rejectMicroserviceMetadataCleanup(format string, args ...any) error { + return errors.Wrapf(errMicroserviceMetadataCleanupRejected, format, args...) +} + +// scheduleMicroserviceMetadataCleanup starts best-effort cleanup of supported +// API-service metadata after the PD leader starts serving. It does not gate +// leader readiness and is scoped to the current leadership term. +func (s *Server) scheduleMicroserviceMetadataCleanup(ctx context.Context) { + if s.IsKeyspaceGroupEnabled() { + return + } + term, ok := s.captureMicroserviceMetadataCleanupTerm() + if !ok { + log.Warn("cannot capture the PD leadership term for microservice metadata cleanup") + return + } + s.serverLoopWg.Add(1) + go func() { + defer logutil.LogPanic() + defer s.serverLoopWg.Done() + if err := s.runMicroserviceMetadataCleanup(ctx, term); err != nil && ctx.Err() == nil { + log.Warn("microservice metadata cleanup stopped before completion", errs.ZapError(err)) + } + }() +} + +func (s *Server) runMicroserviceMetadataCleanup( + ctx context.Context, + term microserviceMetadataCleanupTerm, +) error { + retryInterval := microserviceMetadataCleanupInitialRetryInterval + for { + err := s.cleanupMicroserviceMetadataInPDMode(ctx, term) + if err == nil { + return nil + } + if goerrors.Is(err, errMicroserviceMetadataCleanupRejected) { + log.Warn("cannot safely clean up microservice metadata in PD mode", + errs.ZapError(err)) + return nil + } + if ctx.Err() != nil { + return ctx.Err() + } + + log.Warn("failed to clean up microservice metadata in PD mode, retry later", + zap.Duration("retry-interval", retryInterval), + errs.ZapError(err)) + retryTimer := time.NewTimer(retryInterval) + select { + case <-ctx.Done(): + retryTimer.Stop() + return ctx.Err() + case <-retryTimer.C: + } + retryInterval = min(retryInterval*2, microserviceMetadataCleanupMaxRetryInterval) + } +} + +func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataCleanupTerm, bool) { + if s.member == nil { + return microserviceMetadataCleanupTerm{}, false + } + leadership := s.member.GetLeadership() + if leadership == nil { + return microserviceMetadataCleanupTerm{}, false + } + lease := leadership.GetLease() + if lease == nil { + return microserviceMetadataCleanupTerm{}, false + } + term := microserviceMetadataCleanupTerm{ + leaderKey: leadership.GetLeaderKey(), + leaderValue: leadership.GetLeaderValue(), + leaseID: lease.GetID(), + } + return term, term.leaderKey != "" && term.leaderValue != "" && term.leaseID != 0 +} + +func (s *Server) cleanupMicroserviceMetadataInPDMode( + ctx context.Context, + term microserviceMetadataCleanupTerm, +) error { + if s.IsKeyspaceGroupEnabled() { + return nil + } + + // The persisted keyspace group contains the stale TSO member addresses that + // block a later switch back to API service mode. Keyspace assignment markers + // and the rest of the keyspace group metadata are intentionally left untouched. + // Read the default group first so normal PD-mode leader campaigns do not scan + // all keyspace groups. Non-default groups only matter when cleanup would mutate + // the default group. + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) + resp, err := etcdutil.EtcdKVGetWithContext(ctx, s.client, groupKey) + if err != nil { + return err + } + if len(resp.Kvs) == 0 { + return nil + } + + groupKV := resp.Kvs[0] + group := &endpoint.KeyspaceGroup{} + if err := json.Unmarshal(groupKV.Value, group); err != nil { + return rejectMicroserviceMetadataCleanup( + "cannot decode TSO keyspace group metadata when cleaning up PD mode microservice metadata: %v", err) + } + if group.ID != constant.DefaultKeyspaceGroupID { + return rejectMicroserviceMetadataCleanup("found TSO keyspace group %d at the default group path when cleaning up PD mode microservice metadata", group.ID) + } + if group.IsSplitting() { + return rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting when cleaning up PD mode microservice metadata") + } + if group.IsMerging() { + return rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging when cleaning up PD mode microservice metadata") + } + if len(group.Members) == 0 { + return nil + } + + resp, err = etcdutil.EtcdKVGetWithContext( + ctx, + s.client, + keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID+1), + clientv3.WithRange(clientv3.GetPrefixRangeEnd(keypath.KeyspaceGroupIDPrefix())), + clientv3.WithLimit(1), + ) + if err != nil { + return err + } + if len(resp.Kvs) > 0 { + nonDefaultGroup := &endpoint.KeyspaceGroup{} + if err := json.Unmarshal(resp.Kvs[0].Value, nonDefaultGroup); err != nil { + return rejectMicroserviceMetadataCleanup( + "cannot decode TSO keyspace group metadata when cleaning up PD mode microservice metadata: %v", err) + } + return rejectMicroserviceMetadataCleanup( + "found non-default TSO keyspace group %d when cleaning up PD mode microservice metadata", nonDefaultGroup.ID) + } + + group.Members = nil + value, err := json.Marshal(group) + if err != nil { + return errs.ErrJSONMarshal.Wrap(err).GenWithStackByCause() + } + + failpoint.InjectCall("beforeMicroserviceMetadataCleanupCommit") + txnResp, err := kv.NewSlowLogTxnWithContext(ctx, s.client). + If( + clientv3.Compare(clientv3.Value(term.leaderKey), "=", term.leaderValue), + clientv3.Compare(clientv3.LeaseValue(term.leaderKey), "=", term.leaseID), + clientv3.Compare(clientv3.ModRevision(groupKey), "=", groupKV.ModRevision), + ). + Then(clientv3.OpPut(groupKey, string(value))). + Commit() + if err != nil { + return errs.ErrEtcdTxnInternal.Wrap(err).GenWithStackByCause() + } + if !txnResp.Succeeded { + return errs.ErrEtcdTxnConflict.FastGenByArgs() + } + log.Info("cleaned up microservice metadata in PD mode", + zap.Bool("cleared-default-keyspace-group-members", true)) + return nil +} diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go new file mode 100644 index 0000000000..6b064d3ab5 --- /dev/null +++ b/server/microservice_cleanup_test.go @@ -0,0 +1,520 @@ +// Copyright 2026 TiKV Project Authors. +// +// 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 server + +import ( + "context" + "strconv" + "testing" + "time" + + "github.com/stretchr/testify/require" + clientv3 "go.etcd.io/etcd/client/v3" + + "github.com/pingcap/failpoint" + "github.com/pingcap/kvproto/pkg/keyspacepb" + + "github.com/tikv/pd/pkg/errs" + "github.com/tikv/pd/pkg/keyspace" + "github.com/tikv/pd/pkg/keyspace/constant" + mcs "github.com/tikv/pd/pkg/mcs/utils/constant" + "github.com/tikv/pd/pkg/member" + "github.com/tikv/pd/pkg/storage" + "github.com/tikv/pd/pkg/storage/endpoint" + "github.com/tikv/pd/pkg/storage/kv" + "github.com/tikv/pd/pkg/utils/etcdutil" + "github.com/tikv/pd/pkg/utils/keypath" +) + +func TestCleanupMicroserviceMetadataInPDMode(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12345) + defer keypath.ResetClusterID() + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + staleMembers := []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }} + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + if err := store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: staleMembers, + Keyspaces: []uint32{1}, + }); err != nil { + return err + } + return store.SaveKeyspaceMeta(txn, &keyspacepb.KeyspaceMeta{ + Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: 1}, + Name: "keyspace-1", + Config: map[string]string{ + keyspace.TSOKeyspaceGroupIDKey: strconv.FormatUint(uint64(constant.DefaultKeyspaceGroupID), 10), + "gc_life_time": "10m", + }, + }) + })) + + registryPath := keypath.RegistryPath(mcs.TSOServiceName, "127.0.0.1:3379") + electionPath := keypath.ElectionPath(&keypath.MsParam{ + ServiceName: mcs.TSOServiceName, + GroupID: constant.DefaultKeyspaceGroupID, + }) + timestampPath := keypath.TimestampPath(constant.DefaultKeyspaceGroupID) + for path, value := range map[string]string{ + registryPath: "tso", + electionPath: "primary", + timestampPath: "timestamp", + } { + _, err := client.Put(ctx, path, value) + re.NoError(err) + } + + re.NoError(svr.cleanupMicroserviceMetadataInPDMode(ctx, term)) + re.NoError(svr.cleanupMicroserviceMetadataInPDMode(ctx, term)) + + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 0) + re.NoError(err) + re.Len(groups, 1) + re.Equal(constant.DefaultKeyspaceGroupID, groups[0].ID) + re.Equal(endpoint.Basic.String(), groups[0].UserKind) + re.Empty(groups[0].Members) + re.Equal([]uint32{1}, groups[0].Keyspaces) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + meta, err := store.LoadKeyspaceMeta(txn, 1) + re.NoError(err) + re.NotNil(meta) + re.Equal("0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) + re.Equal("10m", meta.GetConfig()["gc_life_time"]) + return nil + })) + for path, value := range map[string]string{ + registryPath: "tso", + electionPath: "primary", + timestampPath: "timestamp", + } { + resp, err := etcdutil.EtcdKVGet(client, path) + re.NoError(err) + re.Len(resp.Kvs, 1) + re.Equal(value, string(resp.Kvs[0].Value)) + } +} + +func TestScheduleMicroserviceMetadataCleanupDoesNotBlockServing(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12348) + defer keypath.ResetClusterID() + + store := storage.NewStorageWithEtcdBackend(client) + svr, _ := newMicroserviceMetadataCleanupTestServer(t, client, store) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}, + }) + })) + cleanupReached, releaseCleanup := enableMicroserviceMetadataCleanupCommitBlocker(t) + + svr.scheduleMicroserviceMetadataCleanup(ctx) + waitMicroserviceMetadataCleanupCommit(t, cleanupReached) + re.True(svr.member.IsServing()) + releaseCleanup() + re.Eventually(func() bool { + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 1) + return err == nil && len(groups) == 1 && len(groups[0].Members) == 0 + }, 10*time.Second, 100*time.Millisecond) +} + +func TestRunMicroserviceMetadataCleanupStopsOnRejectedState(t *testing.T) { + re := require.New(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12353) + defer keypath.ResetClusterID() + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + if err := store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}, + }); err != nil { + return err + } + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: 1, + UserKind: endpoint.Standard.String(), + }) + })) + + resultCh := make(chan error, 1) + go func() { + resultCh <- svr.runMicroserviceMetadataCleanup(ctx, term) + }() + select { + case err := <-resultCh: + re.NoError(err) + case <-time.After(10 * time.Second): + cancel() + <-resultCh + t.Fatal("cleanup kept retrying a rejected metadata state") + } +} + +func TestMicroserviceMetadataCleanupSkipsNonDefaultGroupsWithoutStaleMembers(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12357) + t.Cleanup(keypath.ResetClusterID) + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: 1, + UserKind: endpoint.Standard.String(), + }) + })) + + re.NoError(svr.cleanupMicroserviceMetadataInPDMode(ctx, term)) + + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + }) + })) + re.NoError(svr.cleanupMicroserviceMetadataInPDMode(ctx, term)) +} + +func TestRunMicroserviceMetadataCleanupStopsOnContextCancellation(t *testing.T) { + re := require.New(t) + ctx, cancel := context.WithCancel(context.Background()) + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12356) + t.Cleanup(keypath.ResetClusterID) + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}, + }) + })) + cancel() + + resultCh := make(chan error, 1) + go func() { + resultCh <- svr.runMicroserviceMetadataCleanup(ctx, term) + }() + select { + case err := <-resultCh: + re.ErrorIs(err, context.Canceled) + case <-time.After(10 * time.Second): + t.Fatal("cleanup did not stop after its context was canceled") + } +} + +func TestCleanupMicroserviceMetadataInPDModeRejectsNonDefaultGroup(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12346) + defer keypath.ResetClusterID() + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + staleMembers := []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }} + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + if err := store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: staleMembers, + }); err != nil { + return err + } + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: 1, + UserKind: endpoint.Standard.String(), + }) + })) + + err := svr.cleanupMicroserviceMetadataInPDMode(ctx, term) + re.ErrorContains(err, "non-default TSO keyspace group 1") + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 0) + re.NoError(err) + re.Len(groups, 2) + re.Equal(staleMembers, groups[0].Members) +} + +func TestCleanupMicroserviceMetadataInPDModeRejectsGroupTransition(t *testing.T) { + testCases := []struct { + name string + transition func(*endpoint.KeyspaceGroup) + errText string + }{ + { + name: "splitting", + transition: func(group *endpoint.KeyspaceGroup) { + group.SplitState = &endpoint.SplitState{SplitSource: group.ID} + }, + errText: "splitting", + }, + { + name: "merging", + transition: func(group *endpoint.KeyspaceGroup) { + group.MergeState = &endpoint.MergeState{MergeList: []uint32{1}} + }, + errText: "merging", + }, + } + + for i, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(uint64(12400 + i)) + defer keypath.ResetClusterID() + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + group := &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}, + } + testCase.transition(group) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, group) + })) + + err := svr.cleanupMicroserviceMetadataInPDMode(ctx, term) + re.ErrorContains(err, testCase.errText) + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 0) + re.NoError(err) + re.Len(groups, 1) + re.Equal(group, groups[0]) + }) + } +} + +func TestCleanupMicroserviceMetadataIsFencedByLeadershipTerm(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12354) + t.Cleanup(keypath.ResetClusterID) + + store := storage.NewStorageWithEtcdBackend(client) + svr, oldTerm := newMicroserviceMetadataCleanupTestServer(t, client, store) + staleMembers := []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }} + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: staleMembers, + Keyspaces: []uint32{1}, + }) + })) + + cleanupReached, releaseCleanup := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := make(chan error, 1) + go func() { + resultCh <- svr.cleanupMicroserviceMetadataInPDMode(ctx, oldTerm) + }() + waitMicroserviceMetadataCleanupCommit(t, cleanupReached) + + // Campaign again with the same member value to prove that comparing only the + // leader value would allow an old-term cleanup to pass. + svr.member.Resign() + re.NoError(svr.member.GetLeadership().Campaign( + testMicroserviceMetadataCleanupLeaseTimeout, + svr.member.MemberValue(), + )) + svr.member.PromoteSelf() + newTerm, ok := svr.captureMicroserviceMetadataCleanupTerm() + re.True(ok) + re.Equal(oldTerm.leaderValue, newTerm.leaderValue) + re.NotEqual(oldTerm.leaseID, newTerm.leaseID) + + releaseCleanup() + re.ErrorIs(waitMicroserviceMetadataCleanupResult(t, resultCh), errs.ErrEtcdTxnConflict) + + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 0) + re.NoError(err) + re.Len(groups, 1) + re.Equal(staleMembers, groups[0].Members) +} + +func TestCleanupMicroserviceMetadataPreservesConcurrentUpdate(t *testing.T) { + re := require.New(t) + ctx := context.Background() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(12355) + t.Cleanup(keypath.ResetClusterID) + + store := storage.NewStorageWithEtcdBackend(client) + svr, term := newMicroserviceMetadataCleanupTestServer(t, client, store) + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3379", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}, + Keyspaces: []uint32{1}, + }) + })) + + cleanupReached, releaseCleanup := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := make(chan error, 1) + go func() { + resultCh <- svr.cleanupMicroserviceMetadataInPDMode(ctx, term) + }() + waitMicroserviceMetadataCleanupCommit(t, cleanupReached) + + newMembers := []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3380", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }} + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + group, err := store.LoadKeyspaceGroup(txn, constant.DefaultKeyspaceGroupID) + if err != nil { + return err + } + group.Members = newMembers + group.Keyspaces = append(group.Keyspaces, 2) + return store.SaveKeyspaceGroup(txn, group) + })) + + releaseCleanup() + re.ErrorIs(waitMicroserviceMetadataCleanupResult(t, resultCh), errs.ErrEtcdTxnConflict) + + groups, err := store.LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 0) + re.NoError(err) + re.Len(groups, 1) + re.Equal(newMembers, groups[0].Members) + re.Equal([]uint32{1, 2}, groups[0].Keyspaces) +} + +const testMicroserviceMetadataCleanupLeaseTimeout = 60 + +func newMicroserviceMetadataCleanupTestServer( + t *testing.T, + client *clientv3.Client, + store storage.Storage, +) (*Server, microserviceMetadataCleanupTerm) { + t.Helper() + pdMember := member.NewMember(nil, client, 1) + pdMember.InitMemberInfo("http://127.0.0.1:2379", "http://127.0.0.1:2380", "pd-test") + re := require.New(t) + re.NoError(pdMember.GetLeadership().Campaign(testMicroserviceMetadataCleanupLeaseTimeout, pdMember.MemberValue())) + pdMember.PromoteSelf() + t.Cleanup(pdMember.Resign) + + svr := &Server{storage: store, client: client, member: pdMember} + term, ok := svr.captureMicroserviceMetadataCleanupTerm() + re.True(ok) + return svr, term +} + +func enableMicroserviceMetadataCleanupCommitBlocker(t *testing.T) (<-chan struct{}, func()) { + t.Helper() + const name = "github.com/tikv/pd/server/beforeMicroserviceMetadataCleanupCommit" + reached := make(chan struct{}, 1) + release := make(chan struct{}) + unblock := func() { + select { + case <-release: + default: + close(release) + } + } + require.NoError(t, failpoint.EnableCall(name, func() { + select { + case reached <- struct{}{}: + default: + } + <-release + })) + t.Cleanup(func() { + unblock() + require.NoError(t, failpoint.Disable(name)) + }) + return reached, unblock +} + +func waitMicroserviceMetadataCleanupCommit(t *testing.T, reached <-chan struct{}) { + t.Helper() + select { + case <-reached: + case <-time.After(10 * time.Second): + t.Fatal("microservice metadata cleanup did not reach the commit hook") + } +} + +func waitMicroserviceMetadataCleanupResult( + t *testing.T, + resultCh <-chan error, +) error { + t.Helper() + select { + case err := <-resultCh: + return err + case <-time.After(10 * time.Second): + t.Fatal("microservice metadata cleanup did not return") + return nil + } +} diff --git a/server/server.go b/server/server.go index 58e3caa31e..de4a1a7f57 100644 --- a/server/server.go +++ b/server/server.go @@ -2068,6 +2068,7 @@ func (s *Server) campaignLeader() { zap.String("leader-name", s.Name()), zap.Duration("total-cost", totalDuration), zap.Duration("cost", enableLeaderDuration)) + s.scheduleMicroserviceMetadataCleanup(ctx) leaderTicker := time.NewTicker(mcs.LeaderTickInterval) defer leaderTicker.Stop() diff --git a/tests/integrations/mcs/tso/server_test.go b/tests/integrations/mcs/tso/server_test.go index d20fd7642a..7c3bd655fb 100644 --- a/tests/integrations/mcs/tso/server_test.go +++ b/tests/integrations/mcs/tso/server_test.go @@ -20,6 +20,7 @@ import ( "fmt" "io" "net/http" + "slices" "strings" "sync" "testing" @@ -34,17 +35,20 @@ import ( "github.com/pingcap/failpoint" "github.com/pingcap/kvproto/pkg/metapb" + "github.com/pingcap/kvproto/pkg/pdpb" pd "github.com/tikv/pd/client" "github.com/tikv/pd/client/opt" "github.com/tikv/pd/client/pkg/caller" "github.com/tikv/pd/pkg/core" + "github.com/tikv/pd/pkg/keyspace" "github.com/tikv/pd/pkg/keyspace/constant" "github.com/tikv/pd/pkg/mcs/discovery" tso "github.com/tikv/pd/pkg/mcs/tso/server" tsoapi "github.com/tikv/pd/pkg/mcs/tso/server/apis/v1" mcs "github.com/tikv/pd/pkg/mcs/utils/constant" "github.com/tikv/pd/pkg/storage/endpoint" + "github.com/tikv/pd/pkg/storage/kv" "github.com/tikv/pd/pkg/utils/etcdutil" "github.com/tikv/pd/pkg/utils/tempurl" "github.com/tikv/pd/pkg/utils/testutil" @@ -769,6 +773,290 @@ func TestTSOServiceSwitch(t *testing.T) { re.NoError(failpoint.Disable("github.com/tikv/pd/client/servicediscovery/fastUpdateServiceMode")) } +const modeSwitchUserKeyspaceName = "mode_switch_user_ks" +const modeSwitchPDModeKeyspaceName = "mode_switch_pd_ks" + +func TestPDModeSwitchBetweenMicroserviceAndMonolithMultipleTimes(t *testing.T) { + re := require.New(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + re.NoError(failpoint.Enable("github.com/tikv/pd/server/skipKeyspaceRegionCheck", "return")) + t.Cleanup(func() { + re.NoError(failpoint.Disable("github.com/tikv/pd/server/skipKeyspaceRegionCheck")) + }) + + tc, err := tests.NewTestClusterWithKeyspaceGroup(ctx, 1, func(conf *config.Config, _ string) { + conf.Microservice.EnableTSODynamicSwitching = false + conf.Keyspace.WaitRegionSplit = false + }) + re.NoError(err) + defer tc.Destroy() + + re.NoError(tc.RunInitialServers()) + pdServer := tc.GetServer(tc.WaitLeader()) + re.NotNil(pdServer) + re.NoError(pdServer.BootstrapCluster()) + userKeyspace, err := pdServer.GetServer().GetKeyspaceManager().CreateKeyspace(&keyspace.CreateKeyspaceRequest{ + Name: modeSwitchUserKeyspaceName, + }) + re.NoError(err) + userKeyspaceID := userKeyspace.GetId() + cleanupReached := make(chan struct{}) + cleanupRelease := make(chan struct{}) + var cleanupReachedOnce, cleanupReleaseOnce sync.Once + const cleanupBlockerName = "github.com/tikv/pd/server/beforeMicroserviceMetadataCleanupCommit" + re.NoError(failpoint.EnableCall(cleanupBlockerName, func() { + cleanupReachedOnce.Do(func() { + close(cleanupReached) + }) + <-cleanupRelease + })) + defer func() { + cleanupReleaseOnce.Do(func() { + close(cleanupRelease) + }) + re.NoError(failpoint.Disable(cleanupBlockerName)) + }() + + waitDefaultKeyspaceGroupReady(re, pdServer, userKeyspaceID) + tsoCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) + re.NoError(err) + defer func() { + if tsoCluster != nil { + tsoCluster.Destroy() + } + }() + + for _, tsoStartsBeforeCleanup := range []bool{true, false} { + waitDefaultKeyspaceGroupReady(re, pdServer, userKeyspaceID) + waitTSOServiceReady(re, tsoCluster) + waitDefaultKeyspaceGroupMembers(re, pdServer, tsoCluster.GetKeyspaceGroupMember()) + checkMicroserviceTSOAvailable(ctx, re, pdServer) + checkKeyspaceTSOAvailable(ctx, re, pdServer, modeSwitchUserKeyspaceName) + tsoCluster.Destroy() + tsoCluster = nil + seedDefaultKeyspaceGroupMembers(ctx, re, pdServer, []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:1", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }}) + + if tsoStartsBeforeCleanup { + pdServer, err = startPDForServiceMode(ctx, tc, pdServer, nil) + re.NoError(err) + select { + case <-cleanupReached: + case <-time.After(10 * time.Second): + t.Fatal("normal-mode PD did not reach the cleanup commit hook") + } + + re.True(pdServer.WaitLeader(), "normal-mode PD did not become ready while mode cleanup was pending") + checkPDLeaderAPIAvailable(ctx, re, pdServer) + checkPDTSOAvailable(ctx, re, pdServer) + checkKeyspaceTSOAvailable(ctx, re, pdServer, modeSwitchUserKeyspaceName) + } else { + pdServer, err = restartPDForServiceMode(ctx, tc, pdServer, nil) + re.NoError(err) + waitDefaultKeyspaceGroupMembers(re, pdServer, nil) + } + if tsoStartsBeforeCleanup { + // Create a keyspace while PD is running in normal mode. + _, err = pdServer.GetServer().GetKeyspaceManager().CreateKeyspace(&keyspace.CreateKeyspaceRequest{ + Name: modeSwitchPDModeKeyspaceName, + }) + re.NoError(err) + } + + tsoCluster, err = tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) + re.NoError(err) + if tsoStartsBeforeCleanup { + for _, server := range tsoCluster.GetServers() { + resp, err := tests.TestDialClient.Get(server.GetAddr() + tsoapi.APIPathPrefix + "/health") + re.NoError(err) + statusCode := resp.StatusCode + re.NoError(resp.Body.Close()) + re.Equal(http.StatusNotFound, statusCode) + } + cleanupReleaseOnce.Do(func() { + close(cleanupRelease) + }) + } + waitTSOServiceReady(re, tsoCluster) + checkPDTSOAvailable(ctx, re, pdServer) + checkKeyspaceTSOAvailable(ctx, re, pdServer, modeSwitchUserKeyspaceName) + waitDefaultKeyspaceGroupMembers(re, pdServer, nil) + + pdServer, err = restartPDForServiceMode(ctx, tc, pdServer, []string{mcs.PDServiceName}) + re.NoError(err) + if tsoStartsBeforeCleanup { + checkKeyspaceTSOAvailable(ctx, re, pdServer, modeSwitchPDModeKeyspaceName) + } + } + + waitDefaultKeyspaceGroupReady(re, pdServer, userKeyspaceID) + waitTSOServiceReady(re, tsoCluster) + waitDefaultKeyspaceGroupMembers(re, pdServer, tsoCluster.GetKeyspaceGroupMember()) + checkMicroserviceTSOAvailable(ctx, re, pdServer) + checkKeyspaceTSOAvailable(ctx, re, pdServer, modeSwitchUserKeyspaceName) +} + +func restartPDForServiceMode( + ctx context.Context, + tc *tests.TestCluster, + oldServer *tests.TestServer, + services []string, +) (*tests.TestServer, error) { + newServer, err := startPDForServiceMode(ctx, tc, oldServer, services) + if err != nil { + return nil, err + } + if !newServer.WaitLeader() { + return nil, fmt.Errorf("PD server %s did not become ready", newServer.GetConfig().Name) + } + return newServer, nil +} + +func startPDForServiceMode( + ctx context.Context, + tc *tests.TestCluster, + oldServer *tests.TestServer, + services []string, +) (*tests.TestServer, error) { + cfg := oldServer.GetConfig() + serverName := cfg.Name + if err := oldServer.Stop(); err != nil { + return nil, err + } + + newServer, err := tests.NewTestServer(ctx, cfg, services) + if err != nil { + return nil, err + } + tc.GetServers()[serverName] = newServer + if err := newServer.Run(); err != nil { + return nil, err + } + return newServer, nil +} + +func waitDefaultKeyspaceGroupReady( + re *require.Assertions, + pdServer *tests.TestServer, + userKeyspaceID uint32, +) { + testutil.Eventually(re, func() bool { + group, err := loadDefaultKeyspaceGroup(pdServer) + return err == nil && slices.Contains(group.Keyspaces, userKeyspaceID) + }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) +} + +func seedDefaultKeyspaceGroupMembers( + ctx context.Context, + re *require.Assertions, + pdServer *tests.TestServer, + members []endpoint.KeyspaceGroupMember, +) { + store := pdServer.GetServer().GetStorage() + re.NoError(store.RunInTxn(ctx, func(txn kv.Txn) error { + group, err := store.LoadKeyspaceGroup(txn, constant.DefaultKeyspaceGroupID) + if err != nil { + return err + } + group.Members = members + return store.SaveKeyspaceGroup(txn, group) + })) +} + +func waitDefaultKeyspaceGroupMembers( + re *require.Assertions, + pdServer *tests.TestServer, + expected []endpoint.KeyspaceGroupMember, +) { + testutil.Eventually(re, func() bool { + group, err := loadDefaultKeyspaceGroup(pdServer) + if err != nil || len(group.Members) != len(expected) { + return false + } + for _, expectedMember := range expected { + if !slices.ContainsFunc(group.Members, func(member endpoint.KeyspaceGroupMember) bool { + return member.IsAddressEquivalent(expectedMember.Address) + }) { + return false + } + } + return true + }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) +} + +func loadDefaultKeyspaceGroup(pdServer *tests.TestServer) (*endpoint.KeyspaceGroup, error) { + groups, err := pdServer.GetServer().GetStorage().LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 1) + if err != nil { + return nil, err + } + if len(groups) != 1 || groups[0] == nil || groups[0].ID != constant.DefaultKeyspaceGroupID { + return nil, fmt.Errorf("default keyspace group %d does not exist", constant.DefaultKeyspaceGroupID) + } + return groups[0], nil +} + +func waitTSOServiceReady(re *require.Assertions, tsoCluster *tests.TestTSOCluster) { + primary := tsoCluster.WaitForDefaultPrimaryServing(re) + testutil.Eventually(re, func() bool { + resp, err := tests.TestDialClient.Get(primary.GetAddr() + tsoapi.APIPathPrefix + "/health") + if err != nil { + return false + } + defer resp.Body.Close() + return resp.StatusCode == http.StatusOK + }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) +} + +func checkPDTSOAvailable(ctx context.Context, re *require.Assertions, pdServer *tests.TestServer) { + cli := utils.SetupClientWithAPIContext(ctx, re, pd.NewAPIContextV1(), []string{pdServer.GetAddr()}) + defer cli.Close() + physical, logical, err := cli.GetTS(ctx) + re.NoError(err) + re.NotZero(tsoutil.ComposeTS(physical, logical)) +} + +func checkPDLeaderAPIAvailable(ctx context.Context, re *require.Assertions, pdServer *tests.TestServer) { + cli := utils.SetupClientWithAPIContext(ctx, re, pd.NewAPIContextV1(), []string{pdServer.GetAddr()}) + defer cli.Close() + _, err := cli.GetAllStores(ctx) + re.NoError(err) +} + +func checkMicroserviceTSOAvailable(ctx context.Context, re *require.Assertions, pdServer *tests.TestServer) { + addr := strings.TrimPrefix(pdServer.GetAddr(), "http://") + conn, err := grpc.Dial(addr, grpc.WithTransportCredentials(insecure.NewCredentials())) //nolint:staticcheck + re.NoError(err) + defer conn.Close() + + resp, err := pdpb.NewPDClient(conn).GetClusterInfo(ctx, &pdpb.GetClusterInfoRequest{}) + re.NoError(err) + re.Equal([]pdpb.ServiceMode{pdpb.ServiceMode_API_SVC_MODE}, resp.GetServiceModes()) + re.NotEmpty(resp.GetTsoUrls()) + + checkPDTSOAvailable(ctx, re, pdServer) +} + +func checkKeyspaceTSOAvailable( + ctx context.Context, + re *require.Assertions, + pdServer *tests.TestServer, + keyspaceName string, +) { + cli := utils.SetupClientWithAPIContext( + ctx, + re, + pd.NewAPIContextV2(keyspaceName), + []string{pdServer.GetAddr()}, + ) + defer cli.Close() + physical, logical, err := cli.GetTS(ctx) + re.NoError(err) + re.NotZero(tsoutil.ComposeTS(physical, logical)) +} + func checkTSOMonotonic(ctx context.Context, pdClient pd.Client, globalLastTS *uint64, count int) error { for range count { physical, logical, err := pdClient.GetTS(ctx)