From 5ee8d8203195eacd58f7609e323658bc60a54fc0 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Tue, 11 Aug 2026 21:23:40 +0800 Subject: [PATCH 1/7] server: add microservice metadata cleanup API Add a synchronous admin endpoint that clears stale TSO keyspace group member assignments in normal PD mode. Fence the cleanup with the exact PD leadership term and keyspace group revisions while preserving durable topology, assignment markers, and timestamps. Add unit, API, and mode-switch coverage. Signed-off-by: Ryan Leung --- server/api/admin.go | 43 ++ server/api/router.go | 1 + server/microservice_cleanup.go | 175 ++++++++ server/microservice_cleanup_test.go | 472 ++++++++++++++++++++++ tests/integrations/mcs/tso/server_test.go | 216 ++++++++++ tests/server/api/admin_test.go | 69 ++++ 6 files changed, 976 insertions(+) create mode 100644 server/microservice_cleanup.go create mode 100644 server/microservice_cleanup_test.go diff --git a/server/api/admin.go b/server/api/admin.go index f298ee4cb0..81fd93f1c2 100644 --- a/server/api/admin.go +++ b/server/api/admin.go @@ -16,6 +16,7 @@ package api import ( "encoding/json" + "errors" "fmt" "io" "net/http" @@ -44,6 +45,12 @@ type RecoveryStatusResponse struct { Marked bool `json:"marked"` } +// CleanupMicroserviceMetadataResponse represents the result of microservice +// metadata cleanup. +type CleanupMicroserviceMetadataResponse struct { + Changed bool `json:"changed"` +} + func newAdminHandler(svr *server.Server, rd *render.Render) *adminHandler { return &adminHandler{ svr: svr, @@ -51,6 +58,42 @@ func newAdminHandler(svr *server.Server, rd *render.Render) *adminHandler { } } +// CleanupMicroserviceMetadata cleans up stale microservice runtime metadata in +// normal PD mode. +// +// @Tags admin +// @Summary Clean up stale microservice runtime metadata in normal PD mode. +// @Produce json +// @Success 200 {object} CleanupMicroserviceMetadataResponse "Whether any metadata was changed." +// @Failure 409 {string} string "The cleanup is rejected in the current state." +// @Failure 500 {string} string "PD failed to clean up the metadata." +// @Failure 503 {string} string "The cleanup is temporarily unavailable." +// @Router /admin/microservice/metadata/cleanup [post] +func (h *adminHandler) CleanupMicroserviceMetadata(w http.ResponseWriter, r *http.Request) { + // The redirect middleware allows callers to explicitly request follower + // handling. Never let that opt-in turn this mutating operation into a + // follower-local write. + if !h.svr.IsServing() { + h.rd.JSON(w, http.StatusServiceUnavailable, server.ErrMicroserviceMetadataCleanupUnavailable.Error()) + return + } + + changed, err := h.svr.CleanupMicroserviceMetadata(r.Context()) + if err != nil { + switch { + case errors.Is(err, server.ErrMicroserviceMetadataCleanupRejected): + h.rd.JSON(w, http.StatusConflict, err.Error()) + case errors.Is(err, server.ErrMicroserviceMetadataCleanupUnavailable): + h.rd.JSON(w, http.StatusServiceUnavailable, err.Error()) + default: + apiutil.ErrorResp(h.rd, w, err) + } + return + } + + h.rd.JSON(w, http.StatusOK, &CleanupMicroserviceMetadataResponse{Changed: changed}) +} + // DeleteRegionCache removes a specific region from cache. // // @Tags admin diff --git a/server/api/router.go b/server/api/router.go index 20f41967c0..728028a865 100644 --- a/server/api/router.go +++ b/server/api/router.go @@ -312,6 +312,7 @@ func createRouter(prefix string, svr *server.Server) *mux.Router { registerFunc(regionResetRouter, "/admin/cache/region/{id}", adminHandler.DeleteRegionCache, setMethods(http.MethodDelete), setAuditBackend(localLog, prometheus)) registerFunc(clusterRouter, "/admin/storage/region/{id}", adminHandler.DeleteRegionStorage, setMethods(http.MethodDelete), setAuditBackend(localLog, prometheus)) registerFunc(regionResetRouter, "/admin/cache/regions", adminHandler.DeleteAllRegionCache, setMethods(http.MethodDelete), setAuditBackend(localLog, prometheus)) + registerFunc(apiRouter, "/admin/microservice/metadata/cleanup", adminHandler.CleanupMicroserviceMetadata, setMethods(http.MethodPost), setAuditBackend(localLog, prometheus)) registerFunc(apiRouter, "/admin/persist-file/{file_name}", adminHandler.SavePersistFile, setMethods(http.MethodPost), setAuditBackend(localLog, prometheus)) registerFunc(apiRouter, "/admin/cluster/markers/snapshot-recovering", adminHandler.isSnapshotRecovering, setMethods(http.MethodGet), setAuditBackend(localLog, prometheus)) registerFunc(apiRouter, "/admin/cluster/markers/snapshot-recovering", adminHandler.markSnapshotRecovering, setMethods(http.MethodPost), setAuditBackend(localLog, prometheus)) diff --git a/server/microservice_cleanup.go b/server/microservice_cleanup.go new file mode 100644 index 0000000000..c6b68cb8b6 --- /dev/null +++ b/server/microservice_cleanup.go @@ -0,0 +1,175 @@ +// 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" + + clientv3 "go.etcd.io/etcd/client/v3" + + "github.com/pingcap/errors" + "github.com/pingcap/failpoint" + + "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" +) + +var ( + // ErrMicroserviceMetadataCleanupRejected indicates that the persisted + // keyspace-group state cannot be safely cleaned up. + ErrMicroserviceMetadataCleanupRejected = errors.New("microservice metadata cleanup rejected") + // ErrMicroserviceMetadataCleanupUnavailable indicates that the request could + // not be completed under the current normal-PD leadership term. + ErrMicroserviceMetadataCleanupUnavailable = errors.New("microservice metadata cleanup unavailable") +) + +type microserviceMetadataCleanupTerm struct { + leaderKey string + leaderValue string + leaseID clientv3.LeaseID +} + +func rejectMicroserviceMetadataCleanup(format string, args ...any) error { + return errors.Wrapf(ErrMicroserviceMetadataCleanupRejected, format, args...) +} + +func unavailableMicroserviceMetadataCleanup(format string, args ...any) error { + return errors.Wrapf(ErrMicroserviceMetadataCleanupUnavailable, format, args...) +} + +// CleanupMicroserviceMetadata clears the persisted Members field of the default +// TSO keyspace group in normal PD mode and reports whether it changed. A nil +// error is fenced by one exact normal-PD leadership term and represents a +// linearizable check that no non-default keyspace group existed at that point. +func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) { + term, err := s.captureMicroserviceMetadataCleanupTerm() + if err != nil { + return false, err + } + if s.IsKeyspaceGroupEnabled() { + return false, rejectMicroserviceMetadataCleanup("pd is already running in microservice mode") + } + + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) + groupPrefixEnd := clientv3.GetPrefixRangeEnd(keypath.KeyspaceGroupIDPrefix()) + operationCtx, cancel := context.WithTimeout(ctx, etcdutil.DefaultRequestTimeout) + defer cancel() + resp, err := s.client.Get( + operationCtx, + groupKey, + clientv3.WithRange(groupPrefixEnd), + clientv3.WithLimit(2), + ) + if err != nil { + return false, unavailableMicroserviceMetadataCleanup( + "failed to read keyspace-group metadata: %v", err) + } + + var ( + group *endpoint.KeyspaceGroup + groupRevision clientv3.Cmp + changed bool + ) + switch { + case len(resp.Kvs) == 0: + groupRevision = clientv3.Compare(clientv3.CreateRevision(groupKey), "=", 0) + case string(resp.Kvs[0].Key) != groupKey || len(resp.Kvs) > 1: + return false, rejectMicroserviceMetadataCleanup("found a non-default TSO keyspace group") + default: + group = &endpoint.KeyspaceGroup{} + if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + if group.ID != constant.DefaultKeyspaceGroupID { + return false, rejectMicroserviceMetadataCleanup( + "found TSO keyspace group %d at the default group path", group.ID) + } + if group.IsSplitting() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting") + } + if group.IsMerging() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging") + } + groupRevision = clientv3.Compare( + clientv3.ModRevision(groupKey), + "=", + resp.Kvs[0].ModRevision, + ) + changed = len(group.Members) > 0 + } + + nonDefaultGroupStart := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID + 1) + comparisons := []clientv3.Cmp{ + clientv3.Compare(clientv3.Value(term.leaderKey), "=", term.leaderValue), + clientv3.Compare(clientv3.LeaseValue(term.leaderKey), "=", term.leaseID), + groupRevision, + clientv3.Compare( + clientv3.CreateRevision(nonDefaultGroupStart).WithRange(groupPrefixEnd), + "=", + 0, + ), + } + operation := clientv3.OpGet(groupKey) + if changed { + group.Members = nil + value, err := json.Marshal(group) + if err != nil { + return false, errs.ErrJSONMarshal.Wrap(err).GenWithStackByCause() + } + operation = clientv3.OpPut(groupKey, string(value)) + } + + failpoint.InjectCall("beforeCleanupMicroserviceMetadataCommit") + txnResp, err := kv.NewSlowLogTxnWithContext(operationCtx, s.client). + If(comparisons...). + Then(operation). + Commit() + if err != nil { + return false, unavailableMicroserviceMetadataCleanup( + "failed to commit microservice metadata cleanup: %v", err) + } + if txnResp.Succeeded { + return changed, nil + } + return false, unavailableMicroserviceMetadataCleanup( + "leadership or keyspace-group metadata changed during cleanup") +} + +func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataCleanupTerm, error) { + if s.client == nil || s.member == nil || !s.member.IsServing() { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leader is not serving") + } + leadership := s.member.GetLeadership() + if leadership == nil || leadership.GetLease() == nil { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leadership is not initialized") + } + term := microserviceMetadataCleanupTerm{ + leaderKey: leadership.GetLeaderKey(), + leaderValue: leadership.GetLeaderValue(), + leaseID: leadership.GetLease().GetID(), + } + if term.leaderKey == "" || term.leaderValue == "" || term.leaseID == 0 { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leadership term is incomplete") + } + return term, nil +} diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go new file mode 100644 index 0000000000..1394dfb532 --- /dev/null +++ b/server/microservice_cleanup_test.go @@ -0,0 +1,472 @@ +// 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" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/pingcap/failpoint" + "github.com/pingcap/kvproto/pkg/keyspacepb" + + "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 TestCleanupMicroserviceMetadataPreservesDefaultGroup(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13001) + ctx := context.Background() + staleMembers := []endpoint.KeyspaceGroupMember{ + {Address: "http://127.0.0.1:3379", Priority: mcs.DefaultKeyspaceGroupReplicaPriority}, + {Address: "http://127.0.0.1:3380", Priority: mcs.DefaultKeyspaceGroupReplicaPriority}, + } + saveMicroserviceMetadataCleanupTestGroup(t, store, &endpoint.KeyspaceGroup{ + ID: constant.DefaultKeyspaceGroupID, + UserKind: endpoint.Standard.String(), + Members: staleMembers, + Keyspaces: []uint32{1, 2, 3}, + }) + require.NoError(t, store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceMeta(txn, &keyspacepb.KeyspaceMeta{ + Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: 1}, + Name: "keyspace-1", + Config: map[string]string{ + keyspace.TSOKeyspaceGroupIDKey: "0", + "gc_life_time": "10m", + }, + }) + })) + const timestampValue = "group-0-timestamp" + timestampPath := keypath.TimestampPath(constant.DefaultKeyspaceGroupID) + _, err := svr.client.Put(ctx, timestampPath, timestampValue) + require.NoError(t, err) + + changed, err := svr.CleanupMicroserviceMetadata(ctx) + require.NoError(t, err) + require.True(t, changed) + + group := loadMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID) + require.Equal(t, constant.DefaultKeyspaceGroupID, group.ID) + require.Equal(t, endpoint.Standard.String(), group.UserKind) + require.Empty(t, group.Members) + require.Equal(t, []uint32{1, 2, 3}, group.Keyspaces) + require.Nil(t, group.SplitState) + require.Nil(t, group.MergeState) + + changed, err = svr.CleanupMicroserviceMetadata(ctx) + require.NoError(t, err) + require.False(t, changed) + + require.NoError(t, store.RunInTxn(ctx, func(txn kv.Txn) error { + meta, err := store.LoadKeyspaceMeta(txn, 1) + require.NoError(t, err) + require.NotNil(t, meta) + require.Equal(t, "0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) + require.Equal(t, "10m", meta.GetConfig()["gc_life_time"]) + return nil + })) + resp, err := etcdutil.EtcdKVGet(svr.client, timestampPath) + require.NoError(t, err) + require.Len(t, resp.Kvs, 1) + require.Equal(t, timestampValue, string(resp.Kvs[0].Value)) +} + +func TestCleanupMicroserviceMetadataRejectsUnsafeState(t *testing.T) { + t.Run("non-default-group", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13002) + defaultGroup := newMicroserviceMetadataCleanupTestDefaultGroup() + saveMicroserviceMetadataCleanupTestGroup(t, store, defaultGroup) + saveMicroserviceMetadataCleanupTestGroup(t, store, &endpoint.KeyspaceGroup{ + ID: 1, + UserKind: endpoint.Basic.String(), + }) + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + require.Equal(t, defaultGroup, loadMicroserviceMetadataCleanupTestGroup( + t, store, constant.DefaultKeyspaceGroupID)) + }) + + testCases := []struct { + name string + clusterID uint64 + transition func(*endpoint.KeyspaceGroup) + }{ + { + name: "splitting", + clusterID: 13003, + transition: func(group *endpoint.KeyspaceGroup) { + group.SplitState = &endpoint.SplitState{SplitSource: group.ID} + }, + }, + { + name: "merging", + clusterID: 13004, + transition: func(group *endpoint.KeyspaceGroup) { + group.MergeState = &endpoint.MergeState{MergeList: []uint32{1}} + }, + }, + } + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, testCase.clusterID) + group := newMicroserviceMetadataCleanupTestDefaultGroup() + testCase.transition(group) + saveMicroserviceMetadataCleanupTestGroup(t, store, group) + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + require.Equal(t, group, loadMicroserviceMetadataCleanupTestGroup( + t, store, constant.DefaultKeyspaceGroupID)) + }) + } +} + +func TestCleanupMicroserviceMetadataRequiresNormalPDLeader(t *testing.T) { + t.Run("microservice-mode", func(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13005) + svr.isKeyspaceGroupEnabled = true + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + }) + + t.Run("not-serving", func(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13006) + svr.member.Resign() + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupUnavailable) + }) +} + +func TestCleanupMicroserviceMetadataIsFencedByExactLeadershipTerm(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13007) + saveMicroserviceMetadataCleanupTestGroup(t, store, newMicroserviceMetadataCleanupTestDefaultGroup()) + oldTerm, err := svr.captureMicroserviceMetadataCleanupTerm() + require.NoError(t, err) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + + // Re-campaign with the same member value. The new lease is what distinguishes + // this leadership term from the one captured by the blocked request. + svr.member.Resign() + require.NoError(t, svr.member.GetLeadership().Campaign( + testMicroserviceMetadataCleanupLeaseTimeout, + svr.member.MemberValue(), + )) + svr.member.PromoteSelf() + newTerm, err := svr.captureMicroserviceMetadataCleanupTerm() + require.NoError(t, err) + require.Equal(t, oldTerm.leaderValue, newTerm.leaderValue) + require.NotEqual(t, oldTerm.leaseID, newTerm.leaseID) + + blocker.releaseCleanup() + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.NotEmpty(t, loadMicroserviceMetadataCleanupTestGroup( + t, store, constant.DefaultKeyspaceGroupID).Members) +} + +func TestCleanupMicroserviceMetadataPreservesConcurrentGroupUpdate(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13008) + saveMicroserviceMetadataCleanupTestGroup(t, store, newMicroserviceMetadataCleanupTestDefaultGroup()) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + + newMembers := []endpoint.KeyspaceGroupMember{{ + Address: "http://127.0.0.1:3380", + Priority: mcs.DefaultKeyspaceGroupReplicaPriority, + }} + updateMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID, func(group *endpoint.KeyspaceGroup) { + group.Members = newMembers + group.Keyspaces = append(group.Keyspaces, 2) + }) + + blocker.releaseCleanup() + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + group := loadMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID) + require.Equal(t, newMembers, group.Members) + require.Equal(t, []uint32{1, 2}, group.Keyspaces) +} + +func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { + t.Run("absent", func(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13009) + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.NoError(t, err) + require.False(t, changed) + }) + + t.Run("empty-members-leadership-change", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13010) + group := newMicroserviceMetadataCleanupTestDefaultGroup() + group.Members = nil + saveMicroserviceMetadataCleanupTestGroup(t, store, group) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + svr.member.Resign() + blocker.releaseCleanup() + + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + }) + + t.Run("empty-members-group-update", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13011) + group := newMicroserviceMetadataCleanupTestDefaultGroup() + group.Members = nil + saveMicroserviceMetadataCleanupTestGroup(t, store, group) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + updateMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID, func(group *endpoint.KeyspaceGroup) { + group.Keyspaces = append(group.Keyspaces, 2) + }) + blocker.releaseCleanup() + + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.Equal(t, []uint32{1, 2}, loadMicroserviceMetadataCleanupTestGroup( + t, store, constant.DefaultKeyspaceGroupID).Keyspaces) + }) + + t.Run("absent-default-created", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13012) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + saveMicroserviceMetadataCleanupTestGroup(t, store, newMicroserviceMetadataCleanupTestDefaultGroup()) + blocker.releaseCleanup() + + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.NotNil(t, loadMicroserviceMetadataCleanupTestGroup( + t, store, constant.DefaultKeyspaceGroupID)) + }) + + t.Run("non-default-created", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13013) + + blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) + resultCh := runMicroserviceMetadataCleanup(t, svr) + blocker.wait(t) + saveMicroserviceMetadataCleanupTestGroup(t, store, &endpoint.KeyspaceGroup{ + ID: 1, + UserKind: endpoint.Basic.String(), + }) + blocker.releaseCleanup() + + result := waitMicroserviceMetadataCleanupResult(t, resultCh) + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.NotNil(t, loadMicroserviceMetadataCleanupTestGroup(t, store, 1)) + }) +} + +const testMicroserviceMetadataCleanupLeaseTimeout = 60 + +type microserviceMetadataCleanupResult struct { + changed bool + err error +} + +type microserviceMetadataCleanupCommitBlocker struct { + name string + reached chan struct{} + release chan struct{} + reachedOnce sync.Once + releaseOnce sync.Once + disableOnce sync.Once +} + +func newMicroserviceMetadataCleanupTestServer( + t *testing.T, + clusterID uint64, +) (*Server, storage.Storage) { + t.Helper() + _, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil) + t.Cleanup(clean) + keypath.SetClusterID(clusterID) + t.Cleanup(keypath.ResetClusterID) + + store := storage.NewStorageWithEtcdBackend(client) + pdMember := member.NewMember(nil, client, 1) + pdMember.InitMemberInfo("http://127.0.0.1:2379", "http://127.0.0.1:2380", "pd-test") + require.NoError(t, pdMember.GetLeadership().Campaign( + testMicroserviceMetadataCleanupLeaseTimeout, + pdMember.MemberValue(), + )) + pdMember.PromoteSelf() + t.Cleanup(pdMember.Resign) + + return &Server{ + storage: store, + client: client, + member: pdMember, + }, store +} + +func newMicroserviceMetadataCleanupTestDefaultGroup() *endpoint.KeyspaceGroup { + return &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}, + } +} + +func saveMicroserviceMetadataCleanupTestGroup( + t *testing.T, + store storage.Storage, + group *endpoint.KeyspaceGroup, +) { + t.Helper() + require.NoError(t, store.RunInTxn(context.Background(), func(txn kv.Txn) error { + return store.SaveKeyspaceGroup(txn, group) + })) +} + +func updateMicroserviceMetadataCleanupTestGroup( + t *testing.T, + store storage.Storage, + id uint32, + update func(*endpoint.KeyspaceGroup), +) { + t.Helper() + require.NoError(t, store.RunInTxn(context.Background(), func(txn kv.Txn) error { + group, err := store.LoadKeyspaceGroup(txn, id) + if err != nil { + return err + } + require.NotNil(t, group) + update(group) + return store.SaveKeyspaceGroup(txn, group) + })) +} + +func loadMicroserviceMetadataCleanupTestGroup( + t *testing.T, + store storage.Storage, + id uint32, +) *endpoint.KeyspaceGroup { + t.Helper() + var group *endpoint.KeyspaceGroup + require.NoError(t, store.RunInTxn(context.Background(), func(txn kv.Txn) error { + var err error + group, err = store.LoadKeyspaceGroup(txn, id) + return err + })) + return group +} + +func runMicroserviceMetadataCleanup(t *testing.T, svr *Server) <-chan microserviceMetadataCleanupResult { + t.Helper() + resultCh := make(chan microserviceMetadataCleanupResult, 1) + go func() { + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + resultCh <- microserviceMetadataCleanupResult{changed: changed, err: err} + }() + return resultCh +} + +func enableMicroserviceMetadataCleanupCommitBlocker(t *testing.T) *microserviceMetadataCleanupCommitBlocker { + t.Helper() + blocker := µserviceMetadataCleanupCommitBlocker{ + name: "github.com/tikv/pd/server/beforeCleanupMicroserviceMetadataCommit", + reached: make(chan struct{}), + release: make(chan struct{}), + } + require.NoError(t, failpoint.EnableCall(blocker.name, func() { + blocker.reachedOnce.Do(func() { + close(blocker.reached) + }) + <-blocker.release + })) + t.Cleanup(func() { + blocker.releaseCleanup() + blocker.disable(t) + }) + return blocker +} + +func (b *microserviceMetadataCleanupCommitBlocker) wait(t *testing.T) { + t.Helper() + select { + case <-b.reached: + case <-time.After(10 * time.Second): + t.Fatal("microservice metadata cleanup did not reach the commit hook") + } +} + +func (b *microserviceMetadataCleanupCommitBlocker) releaseCleanup() { + b.releaseOnce.Do(func() { + close(b.release) + }) +} + +func (b *microserviceMetadataCleanupCommitBlocker) disable(t *testing.T) { + t.Helper() + b.disableOnce.Do(func() { + require.NoError(t, failpoint.Disable(b.name)) + }) +} + +func waitMicroserviceMetadataCleanupResult( + t *testing.T, + resultCh <-chan microserviceMetadataCleanupResult, +) microserviceMetadataCleanupResult { + t.Helper() + select { + case result := <-resultCh: + return result + case <-time.After(10 * time.Second): + t.Fatal("microservice metadata cleanup did not return") + return microserviceMetadataCleanupResult{} + } +} diff --git a/tests/integrations/mcs/tso/server_test.go b/tests/integrations/mcs/tso/server_test.go index d20fd7642a..0e16b80548 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" @@ -39,6 +40,7 @@ import ( "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" @@ -769,6 +771,220 @@ func TestTSOServiceSwitch(t *testing.T) { re.NoError(failpoint.Disable("github.com/tikv/pd/client/servicediscovery/fastUpdateServiceMode")) } +func TestCleanupMicroserviceMetadataForModeSwitch(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")) + }) + re.NoError(failpoint.Enable("github.com/tikv/pd/client/servicediscovery/fastUpdateServiceMode", "return(true)")) + t.Cleanup(func() { + re.NoError(failpoint.Disable("github.com/tikv/pd/client/servicediscovery/fastUpdateServiceMode")) + }) + + 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()) + + const userKeyspaceName = "metadata_cleanup_user_keyspace" + userKeyspace, err := pdServer.GetServer().GetKeyspaceManager().CreateKeyspace(&keyspace.CreateKeyspaceRequest{ + Name: userKeyspaceName, + }) + re.NoError(err) + userKeyspaceID := userKeyspace.GetId() + waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + + oldTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) + re.NoError(err) + defer oldTSOCluster.Destroy() + oldTSOPrimary := waitForTSOServiceReady(re, oldTSOCluster) + oldMembers := oldTSOCluster.GetKeyspaceGroupMember() + waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) + + defaultClient := utils.SetupClientWithAPIContext(ctx, re, pd.NewAPIContextV1(), []string{pdServer.GetAddr()}) + defer defaultClient.Close() + userClient := utils.SetupClientWithAPIContext(ctx, re, pd.NewAPIContextV2(userKeyspaceName), []string{pdServer.GetAddr()}) + defer userClient.Close() + globalLastTS := tsoutil.ComposeTS(time.Now().Add(5*time.Minute).UnixMilli(), 0) + re.NoError(oldTSOPrimary.ResetTS( + globalLastTS, + true, + true, + constant.DefaultKeyspaceGroupID, + )) + re.NoError(checkTSOMonotonic(ctx, defaultClient, &globalLastTS, 10)) + re.NoError(checkTSOMonotonic(ctx, userClient, &globalLastTS, 10)) + + oldTSOCluster.Destroy() + assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + + pdServer = restartPDWithServices(ctx, re, tc, pdServer, nil) + waitForTSOMonotonic(re, ctx, defaultClient, &globalLastTS) + waitForTSOMonotonic(re, ctx, userClient, &globalLastTS) + assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + + re.True(cleanupMicroserviceMetadataViaHTTP(re, pdServer)) + assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + re.NoError(checkTSOMonotonic(ctx, defaultClient, &globalLastTS, 10)) + re.NoError(checkTSOMonotonic(ctx, userClient, &globalLastTS, 10)) + + pdServer = restartPDWithServices(ctx, re, tc, pdServer, []string{mcs.PDServiceName}) + assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + + newTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) + re.NoError(err) + defer newTSOCluster.Destroy() + waitForTSOServiceReady(re, newTSOCluster) + newMembers := newTSOCluster.GetKeyspaceGroupMember() + waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, newMembers) + waitForTSOMonotonic(re, ctx, defaultClient, &globalLastTS) + waitForTSOMonotonic(re, ctx, userClient, &globalLastTS) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) +} + +func restartPDWithServices( + ctx context.Context, + re *require.Assertions, + tc *tests.TestCluster, + oldServer *tests.TestServer, + services []string, +) *tests.TestServer { + cfg := oldServer.GetConfig() + serverName := cfg.Name + re.NoError(oldServer.Stop()) + + newServer, err := tests.NewTestServer(ctx, cfg, services) + re.NoError(err) + tc.GetServers()[serverName] = newServer + re.NoError(newServer.Run()) + re.True(newServer.WaitLeader()) + return newServer +} + +func cleanupMicroserviceMetadataViaHTTP(re *require.Assertions, pdServer *tests.TestServer) bool { + const cleanupPath = "/pd/api/v1/admin/microservice/metadata/cleanup" + req, err := http.NewRequest(http.MethodPost, pdServer.GetAddr()+cleanupPath, http.NoBody) + re.NoError(err) + resp, err := tests.TestDialClient.Do(req) + re.NoError(err) + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + re.NoError(err) + re.Equal(http.StatusOK, resp.StatusCode, string(body)) + + result := struct { + Changed bool `json:"changed"` + }{} + re.NoError(json.Unmarshal(body, &result)) + return result.Changed +} + +func waitForDefaultKeyspaceGroup( + re *require.Assertions, + pdServer *tests.TestServer, + userKeyspaceID uint32, + expectedMembers []endpoint.KeyspaceGroupMember, +) { + testutil.Eventually(re, func() bool { + group, err := loadOnlyDefaultKeyspaceGroup(pdServer) + return err == nil && + slices.Contains(group.Keyspaces, constant.DefaultKeyspaceID) && + slices.Contains(group.Keyspaces, userKeyspaceID) && + keyspaceGroupMembersEqual(group.Members, expectedMembers) + }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) +} + +func assertDefaultKeyspaceGroup( + re *require.Assertions, + pdServer *tests.TestServer, + userKeyspaceID uint32, + expectedMembers []endpoint.KeyspaceGroupMember, +) { + group, err := loadOnlyDefaultKeyspaceGroup(pdServer) + re.NoError(err) + re.Contains(group.Keyspaces, constant.DefaultKeyspaceID) + re.Contains(group.Keyspaces, userKeyspaceID) + re.True(keyspaceGroupMembersEqual(group.Members, expectedMembers), + "unexpected keyspace group members, expected %v, got %v", expectedMembers, group.Members) +} + +func loadOnlyDefaultKeyspaceGroup(pdServer *tests.TestServer) (*endpoint.KeyspaceGroup, error) { + groups, err := pdServer.GetServer().GetStorage().LoadKeyspaceGroups(constant.DefaultKeyspaceGroupID, 2) + if err != nil { + return nil, err + } + if len(groups) != 1 || groups[0] == nil || groups[0].ID != constant.DefaultKeyspaceGroupID { + return nil, fmt.Errorf("expected only default keyspace group %d, got %v", constant.DefaultKeyspaceGroupID, groups) + } + return groups[0], nil +} + +func keyspaceGroupMembersEqual(actual, expected []endpoint.KeyspaceGroupMember) bool { + if len(actual) != len(expected) { + return false + } + for _, expectedMember := range expected { + if !slices.ContainsFunc(actual, func(member endpoint.KeyspaceGroupMember) bool { + return member.IsAddressEquivalent(expectedMember.Address) && member.Priority == expectedMember.Priority + }) { + return false + } + } + return true +} + +func assertUserKeyspaceInDefaultGroup( + re *require.Assertions, + pdServer *tests.TestServer, + userKeyspaceName string, +) { + meta, err := pdServer.GetServer().GetKeyspaceManager().LoadKeyspace(userKeyspaceName) + re.NoError(err) + re.Equal("0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) +} + +func waitForTSOServiceReady(re *require.Assertions, tsoCluster *tests.TestTSOCluster) *tso.Server { + 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)) + return primary +} + +func waitForTSOMonotonic(re *require.Assertions, ctx context.Context, client pd.Client, globalLastTS *uint64) { + for range 10 { + var physical, logical int64 + testutil.Eventually(re, func() bool { + var err error + physical, logical, err = client.GetTS(ctx) + return err == nil + }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) + ts := (uint64(physical) << 18) + uint64(logical) + re.Greater(ts, *globalLastTS) + *globalLastTS = ts + } +} + func checkTSOMonotonic(ctx context.Context, pdClient pd.Client, globalLastTS *uint64, count int) error { for range count { physical, logical, err := pdClient.GetTS(ctx) diff --git a/tests/server/api/admin_test.go b/tests/server/api/admin_test.go index 396cd1523a..d6249b5d9b 100644 --- a/tests/server/api/admin_test.go +++ b/tests/server/api/admin_test.go @@ -31,7 +31,11 @@ import ( "github.com/pingcap/kvproto/pkg/pdpb" "github.com/tikv/pd/pkg/core" + keyspaceconstant "github.com/tikv/pd/pkg/keyspace/constant" + mcsconstant "github.com/tikv/pd/pkg/mcs/utils/constant" "github.com/tikv/pd/pkg/replication" + "github.com/tikv/pd/pkg/storage/endpoint" + "github.com/tikv/pd/pkg/storage/kv" "github.com/tikv/pd/pkg/utils/apiutil" "github.com/tikv/pd/pkg/utils/testutil" "github.com/tikv/pd/server" @@ -352,6 +356,71 @@ func (suite *adminTestSuite) checkMarkPitrRestoreMode(cluster *tests.TestCluster testutil.StatusOK(re), testutil.StringContain(re, "false"))) } +func (suite *adminTestSuite) TestCleanupMicroserviceMetadata() { + suite.env.RunTestInNonMicroserviceEnv(suite.checkCleanupMicroserviceMetadata) +} + +func (suite *adminTestSuite) checkCleanupMicroserviceMetadata(cluster *tests.TestCluster) { + re := suite.Require() + leader := cluster.GetLeaderServer() + storage := leader.GetServer().GetStorage() + ctx := context.Background() + + groups, err := storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) + re.NoError(err) + var previousDefaultGroup *endpoint.KeyspaceGroup + if len(groups) > 0 && groups[0].ID == keyspaceconstant.DefaultKeyspaceGroupID { + previousDefaultGroup = groups[0] + } + defer func() { + re.NoError(storage.RunInTxn(ctx, func(txn kv.Txn) error { + if previousDefaultGroup == nil { + return storage.DeleteKeyspaceGroup(txn, keyspaceconstant.DefaultKeyspaceGroupID) + } + return storage.SaveKeyspaceGroup(txn, previousDefaultGroup) + })) + }() + + group := &endpoint.KeyspaceGroup{ + ID: keyspaceconstant.DefaultKeyspaceGroupID, + UserKind: endpoint.Basic.String(), + Members: []endpoint.KeyspaceGroupMember{ + { + Address: "http://127.0.0.1:3379", + Priority: mcsconstant.DefaultKeyspaceGroupReplicaPriority, + }, + }, + Keyspaces: []uint32{keyspaceconstant.DefaultKeyspaceID, 42}, + } + re.NoError(storage.RunInTxn(ctx, func(txn kv.Txn) error { + return storage.SaveKeyspaceGroup(txn, group) + })) + + url := leader.GetAddr() + "/pd/api/v1/admin/microservice/metadata/cleanup" + // Only POST is allowed. A rejected method must not mutate the metadata. + re.NoError(testutil.CheckGetJSON(tests.TestDialClient, url, nil, + testutil.Status(re, http.StatusMethodNotAllowed))) + storedGroups, err := storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) + re.NoError(err) + re.Len(storedGroups, 1) + re.Equal(group.Members, storedGroups[0].Members) + + response := &api.CleanupMicroserviceMetadataResponse{} + re.NoError(testutil.CheckPostJSON(tests.TestDialClient, url, nil, + testutil.StatusOK(re), testutil.ExtractJSON(re, response))) + re.True(response.Changed) + storedGroups, err = storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) + re.NoError(err) + re.Len(storedGroups, 1) + re.Empty(storedGroups[0].Members) + re.Equal(group.Keyspaces, storedGroups[0].Keyspaces) + + response.Changed = true + re.NoError(testutil.CheckPostJSON(tests.TestDialClient, url, nil, + testutil.StatusOK(re), testutil.ExtractJSON(re, response))) + re.False(response.Changed) +} + func (suite *adminTestSuite) TestRecoverAllocID() { suite.env.RunTest(suite.checkRecoverAllocID) } From be4caaded31fe805f2908b8a61a7bf5ff6c2ea6b Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 11:23:30 +0800 Subject: [PATCH 2/7] server: reject cleanup without default group Fail closed when the default TSO keyspace group is missing so the cleanup API cannot certify an ambiguous assignment state. Preserve revision fencing when the group is deleted concurrently. Fix the mode-switch integration test input and static-check issues, and cover the missing-group HTTP behavior. Signed-off-by: Ryan Leung --- server/microservice_cleanup.go | 60 +++++++++++------------ server/microservice_cleanup_test.go | 57 ++++++++++++++++----- tests/integrations/mcs/tso/server_test.go | 35 +++++++------ tests/server/api/admin_test.go | 9 ++++ 4 files changed, 100 insertions(+), 61 deletions(-) diff --git a/server/microservice_cleanup.go b/server/microservice_cleanup.go index c6b68cb8b6..0d5e87065e 100644 --- a/server/microservice_cleanup.go +++ b/server/microservice_cleanup.go @@ -57,7 +57,10 @@ func unavailableMicroserviceMetadataCleanup(format string, args ...any) error { // CleanupMicroserviceMetadata clears the persisted Members field of the default // TSO keyspace group in normal PD mode and reports whether it changed. A nil // error is fenced by one exact normal-PD leadership term and represents a -// linearizable check that no non-default keyspace group existed at that point. +// linearizable check that the default group existed, had no transition in +// progress, had no non-default sibling, and had no persisted member at that +// point. The caller must ensure keyspace assignments have already been merged +// into the default group and all microservice metadata writers are stopped. func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) { term, err := s.captureMicroserviceMetadataCleanupTerm() if err != nil { @@ -82,39 +85,34 @@ func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) "failed to read keyspace-group metadata: %v", err) } - var ( - group *endpoint.KeyspaceGroup - groupRevision clientv3.Cmp - changed bool - ) - switch { - case len(resp.Kvs) == 0: - groupRevision = clientv3.Compare(clientv3.CreateRevision(groupKey), "=", 0) - case string(resp.Kvs[0].Key) != groupKey || len(resp.Kvs) > 1: + if len(resp.Kvs) == 0 { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group does not exist") + } + if string(resp.Kvs[0].Key) != groupKey || len(resp.Kvs) > 1 { return false, rejectMicroserviceMetadataCleanup("found a non-default TSO keyspace group") - default: - group = &endpoint.KeyspaceGroup{} - if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { - return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() - } - if group.ID != constant.DefaultKeyspaceGroupID { - return false, rejectMicroserviceMetadataCleanup( - "found TSO keyspace group %d at the default group path", group.ID) - } - if group.IsSplitting() { - return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting") - } - if group.IsMerging() { - return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging") - } - groupRevision = clientv3.Compare( - clientv3.ModRevision(groupKey), - "=", - resp.Kvs[0].ModRevision, - ) - changed = len(group.Members) > 0 } + group := &endpoint.KeyspaceGroup{} + if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + if group.ID != constant.DefaultKeyspaceGroupID { + return false, rejectMicroserviceMetadataCleanup( + "found TSO keyspace group %d at the default group path", group.ID) + } + if group.IsSplitting() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting") + } + if group.IsMerging() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging") + } + groupRevision := clientv3.Compare( + clientv3.ModRevision(groupKey), + "=", + resp.Kvs[0].ModRevision, + ) + changed := len(group.Members) > 0 + nonDefaultGroupStart := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID + 1) comparisons := []clientv3.Cmp{ clientv3.Compare(clientv3.Value(term.leaderKey), "=", term.leaderValue), diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go index 1394dfb532..5b4c37cdbb 100644 --- a/server/microservice_cleanup_test.go +++ b/server/microservice_cleanup_test.go @@ -95,6 +95,31 @@ func TestCleanupMicroserviceMetadataPreservesDefaultGroup(t *testing.T) { } func TestCleanupMicroserviceMetadataRejectsUnsafeState(t *testing.T) { + t.Run("missing-default-group", func(t *testing.T) { + svr, store := newMicroserviceMetadataCleanupTestServer(t, 13014) + ctx := context.Background() + require.NoError(t, store.RunInTxn(ctx, func(txn kv.Txn) error { + return store.SaveKeyspaceMeta(txn, &keyspacepb.KeyspaceMeta{ + Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: 1}, + Name: "keyspace-1", + Config: map[string]string{ + keyspace.TSOKeyspaceGroupIDKey: "0", + }, + }) + })) + + changed, err := svr.CleanupMicroserviceMetadata(ctx) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + require.NoError(t, store.RunInTxn(ctx, func(txn kv.Txn) error { + meta, err := store.LoadKeyspaceMeta(txn, 1) + require.NoError(t, err) + require.NotNil(t, meta) + require.Equal(t, "0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) + return nil + })) + }) + t.Run("non-default-group", func(t *testing.T) { svr, store := newMicroserviceMetadataCleanupTestServer(t, 13002) defaultGroup := newMicroserviceMetadataCleanupTestDefaultGroup() @@ -194,6 +219,7 @@ func TestCleanupMicroserviceMetadataIsFencedByExactLeadershipTerm(t *testing.T) result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") require.NotEmpty(t, loadMicroserviceMetadataCleanupTestGroup( t, store, constant.DefaultKeyspaceGroupID).Members) } @@ -219,20 +245,13 @@ func TestCleanupMicroserviceMetadataPreservesConcurrentGroupUpdate(t *testing.T) result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") group := loadMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID) require.Equal(t, newMembers, group.Members) require.Equal(t, []uint32{1, 2}, group.Keyspaces) } func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { - t.Run("absent", func(t *testing.T) { - svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13009) - - changed, err := svr.CleanupMicroserviceMetadata(context.Background()) - require.NoError(t, err) - require.False(t, changed) - }) - t.Run("empty-members-leadership-change", func(t *testing.T) { svr, store := newMicroserviceMetadataCleanupTestServer(t, 13010) group := newMicroserviceMetadataCleanupTestDefaultGroup() @@ -248,6 +267,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") }) t.Run("empty-members-group-update", func(t *testing.T) { @@ -267,28 +287,40 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") require.Equal(t, []uint32{1, 2}, loadMicroserviceMetadataCleanupTestGroup( t, store, constant.DefaultKeyspaceGroupID).Keyspaces) }) - t.Run("absent-default-created", func(t *testing.T) { + t.Run("empty-members-group-deleted", func(t *testing.T) { svr, store := newMicroserviceMetadataCleanupTestServer(t, 13012) + group := newMicroserviceMetadataCleanupTestDefaultGroup() + group.Members = nil + saveMicroserviceMetadataCleanupTestGroup(t, store, group) blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) resultCh := runMicroserviceMetadataCleanup(t, svr) blocker.wait(t) - saveMicroserviceMetadataCleanupTestGroup(t, store, newMicroserviceMetadataCleanupTestDefaultGroup()) + require.NoError(t, store.RunInTxn(context.Background(), func(txn kv.Txn) error { + return store.DeleteKeyspaceGroup(txn, constant.DefaultKeyspaceGroupID) + })) blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.NotNil(t, loadMicroserviceMetadataCleanupTestGroup( - t, store, constant.DefaultKeyspaceGroupID)) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + + changed, err := svr.CleanupMicroserviceMetadata(context.Background()) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) }) t.Run("non-default-created", func(t *testing.T) { svr, store := newMicroserviceMetadataCleanupTestServer(t, 13013) + group := newMicroserviceMetadataCleanupTestDefaultGroup() + group.Members = nil + saveMicroserviceMetadataCleanupTestGroup(t, store, group) blocker := enableMicroserviceMetadataCleanupCommitBlocker(t) resultCh := runMicroserviceMetadataCleanup(t, svr) @@ -302,6 +334,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { result := waitMicroserviceMetadataCleanupResult(t, resultCh) require.False(t, result.changed) require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") require.NotNil(t, loadMicroserviceMetadataCleanupTestGroup(t, store, 1)) }) } diff --git a/tests/integrations/mcs/tso/server_test.go b/tests/integrations/mcs/tso/server_test.go index 0e16b80548..2c878a8b5b 100644 --- a/tests/integrations/mcs/tso/server_test.go +++ b/tests/integrations/mcs/tso/server_test.go @@ -52,6 +52,7 @@ import ( "github.com/tikv/pd/pkg/utils/testutil" "github.com/tikv/pd/pkg/utils/tsoutil" "github.com/tikv/pd/pkg/versioninfo/kerneltype" + serverapi "github.com/tikv/pd/server/api" "github.com/tikv/pd/server/config" "github.com/tikv/pd/tests" "github.com/tikv/pd/tests/integrations/mcs/utils" @@ -796,14 +797,14 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { re.NotNil(pdServer) re.NoError(pdServer.BootstrapCluster()) - const userKeyspaceName = "metadata_cleanup_user_keyspace" + const userKeyspaceName = "metadata_cleanup_ks" userKeyspace, err := pdServer.GetServer().GetKeyspaceManager().CreateKeyspace(&keyspace.CreateKeyspaceRequest{ Name: userKeyspaceName, }) re.NoError(err) userKeyspaceID := userKeyspace.GetId() waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) oldTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) re.NoError(err) @@ -828,23 +829,23 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { oldTSOCluster.Destroy() assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) pdServer = restartPDWithServices(ctx, re, tc, pdServer, nil) - waitForTSOMonotonic(re, ctx, defaultClient, &globalLastTS) - waitForTSOMonotonic(re, ctx, userClient, &globalLastTS) + waitForTSOMonotonic(ctx, re, defaultClient, &globalLastTS) + waitForTSOMonotonic(ctx, re, userClient, &globalLastTS) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) re.True(cleanupMicroserviceMetadataViaHTTP(re, pdServer)) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) re.NoError(checkTSOMonotonic(ctx, defaultClient, &globalLastTS, 10)) re.NoError(checkTSOMonotonic(ctx, userClient, &globalLastTS, 10)) pdServer = restartPDWithServices(ctx, re, tc, pdServer, []string{mcs.PDServiceName}) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) newTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) re.NoError(err) @@ -852,9 +853,9 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { waitForTSOServiceReady(re, newTSOCluster) newMembers := newTSOCluster.GetKeyspaceGroupMember() waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, newMembers) - waitForTSOMonotonic(re, ctx, defaultClient, &globalLastTS) - waitForTSOMonotonic(re, ctx, userClient, &globalLastTS) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceName) + waitForTSOMonotonic(ctx, re, defaultClient, &globalLastTS) + waitForTSOMonotonic(ctx, re, userClient, &globalLastTS) + assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) } func restartPDWithServices( @@ -887,9 +888,7 @@ func cleanupMicroserviceMetadataViaHTTP(re *require.Assertions, pdServer *tests. re.NoError(err) re.Equal(http.StatusOK, resp.StatusCode, string(body)) - result := struct { - Changed bool `json:"changed"` - }{} + result := serverapi.CleanupMicroserviceMetadataResponse{} re.NoError(json.Unmarshal(body, &result)) return result.Changed } @@ -951,9 +950,9 @@ func keyspaceGroupMembersEqual(actual, expected []endpoint.KeyspaceGroupMember) func assertUserKeyspaceInDefaultGroup( re *require.Assertions, pdServer *tests.TestServer, - userKeyspaceName string, + userKeyspaceID uint32, ) { - meta, err := pdServer.GetServer().GetKeyspaceManager().LoadKeyspace(userKeyspaceName) + meta, err := pdServer.GetServer().GetKeyspaceManager().LoadKeyspaceByID(userKeyspaceID) re.NoError(err) re.Equal("0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) } @@ -971,7 +970,7 @@ func waitForTSOServiceReady(re *require.Assertions, tsoCluster *tests.TestTSOClu return primary } -func waitForTSOMonotonic(re *require.Assertions, ctx context.Context, client pd.Client, globalLastTS *uint64) { +func waitForTSOMonotonic(ctx context.Context, re *require.Assertions, client pd.Client, globalLastTS *uint64) { for range 10 { var physical, logical int64 testutil.Eventually(re, func() bool { @@ -979,7 +978,7 @@ func waitForTSOMonotonic(re *require.Assertions, ctx context.Context, client pd. physical, logical, err = client.GetTS(ctx) return err == nil }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) - ts := (uint64(physical) << 18) + uint64(logical) + ts := tsoutil.ComposeTS(physical, logical) re.Greater(ts, *globalLastTS) *globalLastTS = ts } diff --git a/tests/server/api/admin_test.go b/tests/server/api/admin_test.go index d6249b5d9b..fd1c225cbd 100644 --- a/tests/server/api/admin_test.go +++ b/tests/server/api/admin_test.go @@ -397,7 +397,16 @@ func (suite *adminTestSuite) checkCleanupMicroserviceMetadata(cluster *tests.Tes })) url := leader.GetAddr() + "/pd/api/v1/admin/microservice/metadata/cleanup" + re.NoError(storage.RunInTxn(ctx, func(txn kv.Txn) error { + return storage.DeleteKeyspaceGroup(txn, keyspaceconstant.DefaultKeyspaceGroupID) + })) + re.NoError(testutil.CheckPostJSON(tests.TestDialClient, url, nil, + testutil.Status(re, http.StatusConflict))) + // Only POST is allowed. A rejected method must not mutate the metadata. + re.NoError(storage.RunInTxn(ctx, func(txn kv.Txn) error { + return storage.SaveKeyspaceGroup(txn, group) + })) re.NoError(testutil.CheckGetJSON(tests.TestDialClient, url, nil, testutil.Status(re, http.StatusMethodNotAllowed))) storedGroups, err := storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) From 5b2a95feb922a26566d867116cf36c95fc61b20a Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 11:54:13 +0800 Subject: [PATCH 3/7] tests: use bootstrap keyspace in cleanup test Signed-off-by: Ryan Leung --- tests/integrations/mcs/tso/server_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/integrations/mcs/tso/server_test.go b/tests/integrations/mcs/tso/server_test.go index 2c878a8b5b..8de7fb7b99 100644 --- a/tests/integrations/mcs/tso/server_test.go +++ b/tests/integrations/mcs/tso/server_test.go @@ -902,7 +902,7 @@ func waitForDefaultKeyspaceGroup( testutil.Eventually(re, func() bool { group, err := loadOnlyDefaultKeyspaceGroup(pdServer) return err == nil && - slices.Contains(group.Keyspaces, constant.DefaultKeyspaceID) && + slices.Contains(group.Keyspaces, keyspace.GetBootstrapKeyspaceID()) && slices.Contains(group.Keyspaces, userKeyspaceID) && keyspaceGroupMembersEqual(group.Members, expectedMembers) }, testutil.WithWaitFor(10*time.Second), testutil.WithTickInterval(100*time.Millisecond)) @@ -916,7 +916,7 @@ func assertDefaultKeyspaceGroup( ) { group, err := loadOnlyDefaultKeyspaceGroup(pdServer) re.NoError(err) - re.Contains(group.Keyspaces, constant.DefaultKeyspaceID) + re.Contains(group.Keyspaces, keyspace.GetBootstrapKeyspaceID()) re.Contains(group.Keyspaces, userKeyspaceID) re.True(keyspaceGroupMembersEqual(group.Members, expectedMembers), "unexpected keyspace group members, expected %v, got %v", expectedMembers, group.Members) From c18b65f451914e22f6f178886400929960d66f8e Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 13:47:16 +0800 Subject: [PATCH 4/7] server: move metadata cleanup into server implementation Signed-off-by: Ryan Leung --- server/microservice_cleanup.go | 173 --------------------------------- server/server.go | 142 +++++++++++++++++++++++++++ 2 files changed, 142 insertions(+), 173 deletions(-) delete mode 100644 server/microservice_cleanup.go diff --git a/server/microservice_cleanup.go b/server/microservice_cleanup.go deleted file mode 100644 index 0d5e87065e..0000000000 --- a/server/microservice_cleanup.go +++ /dev/null @@ -1,173 +0,0 @@ -// 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" - - clientv3 "go.etcd.io/etcd/client/v3" - - "github.com/pingcap/errors" - "github.com/pingcap/failpoint" - - "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" -) - -var ( - // ErrMicroserviceMetadataCleanupRejected indicates that the persisted - // keyspace-group state cannot be safely cleaned up. - ErrMicroserviceMetadataCleanupRejected = errors.New("microservice metadata cleanup rejected") - // ErrMicroserviceMetadataCleanupUnavailable indicates that the request could - // not be completed under the current normal-PD leadership term. - ErrMicroserviceMetadataCleanupUnavailable = errors.New("microservice metadata cleanup unavailable") -) - -type microserviceMetadataCleanupTerm struct { - leaderKey string - leaderValue string - leaseID clientv3.LeaseID -} - -func rejectMicroserviceMetadataCleanup(format string, args ...any) error { - return errors.Wrapf(ErrMicroserviceMetadataCleanupRejected, format, args...) -} - -func unavailableMicroserviceMetadataCleanup(format string, args ...any) error { - return errors.Wrapf(ErrMicroserviceMetadataCleanupUnavailable, format, args...) -} - -// CleanupMicroserviceMetadata clears the persisted Members field of the default -// TSO keyspace group in normal PD mode and reports whether it changed. A nil -// error is fenced by one exact normal-PD leadership term and represents a -// linearizable check that the default group existed, had no transition in -// progress, had no non-default sibling, and had no persisted member at that -// point. The caller must ensure keyspace assignments have already been merged -// into the default group and all microservice metadata writers are stopped. -func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) { - term, err := s.captureMicroserviceMetadataCleanupTerm() - if err != nil { - return false, err - } - if s.IsKeyspaceGroupEnabled() { - return false, rejectMicroserviceMetadataCleanup("pd is already running in microservice mode") - } - - groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) - groupPrefixEnd := clientv3.GetPrefixRangeEnd(keypath.KeyspaceGroupIDPrefix()) - operationCtx, cancel := context.WithTimeout(ctx, etcdutil.DefaultRequestTimeout) - defer cancel() - resp, err := s.client.Get( - operationCtx, - groupKey, - clientv3.WithRange(groupPrefixEnd), - clientv3.WithLimit(2), - ) - if err != nil { - return false, unavailableMicroserviceMetadataCleanup( - "failed to read keyspace-group metadata: %v", err) - } - - if len(resp.Kvs) == 0 { - return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group does not exist") - } - if string(resp.Kvs[0].Key) != groupKey || len(resp.Kvs) > 1 { - return false, rejectMicroserviceMetadataCleanup("found a non-default TSO keyspace group") - } - - group := &endpoint.KeyspaceGroup{} - if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { - return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() - } - if group.ID != constant.DefaultKeyspaceGroupID { - return false, rejectMicroserviceMetadataCleanup( - "found TSO keyspace group %d at the default group path", group.ID) - } - if group.IsSplitting() { - return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting") - } - if group.IsMerging() { - return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging") - } - groupRevision := clientv3.Compare( - clientv3.ModRevision(groupKey), - "=", - resp.Kvs[0].ModRevision, - ) - changed := len(group.Members) > 0 - - nonDefaultGroupStart := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID + 1) - comparisons := []clientv3.Cmp{ - clientv3.Compare(clientv3.Value(term.leaderKey), "=", term.leaderValue), - clientv3.Compare(clientv3.LeaseValue(term.leaderKey), "=", term.leaseID), - groupRevision, - clientv3.Compare( - clientv3.CreateRevision(nonDefaultGroupStart).WithRange(groupPrefixEnd), - "=", - 0, - ), - } - operation := clientv3.OpGet(groupKey) - if changed { - group.Members = nil - value, err := json.Marshal(group) - if err != nil { - return false, errs.ErrJSONMarshal.Wrap(err).GenWithStackByCause() - } - operation = clientv3.OpPut(groupKey, string(value)) - } - - failpoint.InjectCall("beforeCleanupMicroserviceMetadataCommit") - txnResp, err := kv.NewSlowLogTxnWithContext(operationCtx, s.client). - If(comparisons...). - Then(operation). - Commit() - if err != nil { - return false, unavailableMicroserviceMetadataCleanup( - "failed to commit microservice metadata cleanup: %v", err) - } - if txnResp.Succeeded { - return changed, nil - } - return false, unavailableMicroserviceMetadataCleanup( - "leadership or keyspace-group metadata changed during cleanup") -} - -func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataCleanupTerm, error) { - if s.client == nil || s.member == nil || !s.member.IsServing() { - return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leader is not serving") - } - leadership := s.member.GetLeadership() - if leadership == nil || leadership.GetLease() == nil { - return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leadership is not initialized") - } - term := microserviceMetadataCleanupTerm{ - leaderKey: leadership.GetLeaderKey(), - leaderValue: leadership.GetLeaderValue(), - leaseID: leadership.GetLease().GetID(), - } - if term.leaderKey == "" || term.leaderValue == "" || term.leaseID == 0 { - return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leadership term is incomplete") - } - return term, nil -} diff --git a/server/server.go b/server/server.go index 58e3caa31e..303a94db51 100644 --- a/server/server.go +++ b/server/server.go @@ -17,6 +17,7 @@ package server import ( "bytes" "context" + "encoding/json" "fmt" "math" "math/rand/v2" @@ -2248,6 +2249,147 @@ func (s *Server) IsTTLConfigExist(key string) bool { return false } +var ( + // ErrMicroserviceMetadataCleanupRejected indicates that the persisted + // keyspace-group state cannot be safely cleaned up. + ErrMicroserviceMetadataCleanupRejected = errors.New("microservice metadata cleanup rejected") + // ErrMicroserviceMetadataCleanupUnavailable indicates that the request could + // not be completed under the current normal-PD leadership term. + ErrMicroserviceMetadataCleanupUnavailable = errors.New("microservice metadata cleanup unavailable") +) + +type microserviceMetadataCleanupTerm struct { + leaderKey string + leaderValue string + leaseID clientv3.LeaseID +} + +func rejectMicroserviceMetadataCleanup(format string, args ...any) error { + return errors.Wrapf(ErrMicroserviceMetadataCleanupRejected, format, args...) +} + +func unavailableMicroserviceMetadataCleanup(format string, args ...any) error { + return errors.Wrapf(ErrMicroserviceMetadataCleanupUnavailable, format, args...) +} + +// CleanupMicroserviceMetadata clears the persisted Members field of the default +// TSO keyspace group in normal PD mode and reports whether it changed. A nil +// error is fenced by one exact normal-PD leadership term and represents a +// linearizable check that the default group existed, had no transition in +// progress, had no non-default sibling, and had no persisted member at that +// point. The caller must ensure keyspace assignments have already been merged +// into the default group and all microservice metadata writers are stopped. +func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) { + term, err := s.captureMicroserviceMetadataCleanupTerm() + if err != nil { + return false, err + } + if s.IsKeyspaceGroupEnabled() { + return false, rejectMicroserviceMetadataCleanup("pd is already running in microservice mode") + } + + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) + groupPrefixEnd := clientv3.GetPrefixRangeEnd(keypath.KeyspaceGroupIDPrefix()) + operationCtx, cancel := context.WithTimeout(ctx, etcdutil.DefaultRequestTimeout) + defer cancel() + resp, err := s.client.Get( + operationCtx, + groupKey, + clientv3.WithRange(groupPrefixEnd), + clientv3.WithLimit(2), + ) + if err != nil { + return false, unavailableMicroserviceMetadataCleanup( + "failed to read keyspace-group metadata: %v", err) + } + + if len(resp.Kvs) == 0 { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group does not exist") + } + if string(resp.Kvs[0].Key) != groupKey || len(resp.Kvs) > 1 { + return false, rejectMicroserviceMetadataCleanup("found a non-default TSO keyspace group") + } + + group := &endpoint.KeyspaceGroup{} + if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + if group.ID != constant.DefaultKeyspaceGroupID { + return false, rejectMicroserviceMetadataCleanup( + "found TSO keyspace group %d at the default group path", group.ID) + } + if group.IsSplitting() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is splitting") + } + if group.IsMerging() { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group is merging") + } + groupRevision := clientv3.Compare( + clientv3.ModRevision(groupKey), + "=", + resp.Kvs[0].ModRevision, + ) + changed := len(group.Members) > 0 + + nonDefaultGroupStart := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID + 1) + comparisons := []clientv3.Cmp{ + clientv3.Compare(clientv3.Value(term.leaderKey), "=", term.leaderValue), + clientv3.Compare(clientv3.LeaseValue(term.leaderKey), "=", term.leaseID), + groupRevision, + clientv3.Compare( + clientv3.CreateRevision(nonDefaultGroupStart).WithRange(groupPrefixEnd), + "=", + 0, + ), + } + operation := clientv3.OpGet(groupKey) + if changed { + group.Members = nil + value, err := json.Marshal(group) + if err != nil { + return false, errs.ErrJSONMarshal.Wrap(err).GenWithStackByCause() + } + operation = clientv3.OpPut(groupKey, string(value)) + } + + failpoint.InjectCall("beforeCleanupMicroserviceMetadataCommit") + txnResp, err := kv.NewSlowLogTxnWithContext(operationCtx, s.client). + If(comparisons...). + Then(operation). + Commit() + if err != nil { + return false, unavailableMicroserviceMetadataCleanup( + "failed to commit microservice metadata cleanup: %v", err) + } + if txnResp.Succeeded { + return changed, nil + } + return false, unavailableMicroserviceMetadataCleanup( + "leadership or keyspace-group metadata changed during cleanup") +} + +func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataCleanupTerm, error) { + if s.client == nil || s.member == nil || !s.member.IsServing() { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leader is not serving") + } + leadership := s.member.GetLeadership() + if leadership == nil || leadership.GetLease() == nil { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leadership is not initialized") + } + term := microserviceMetadataCleanupTerm{ + leaderKey: leadership.GetLeaderKey(), + leaderValue: leadership.GetLeaderValue(), + leaseID: leadership.GetLease().GetID(), + } + if term.leaderKey == "" || term.leaderValue == "" || term.leaseID == 0 { + return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( + "normal PD leadership term is incomplete") + } + return term, nil +} + // MarkSnapshotRecovering mark pd that we're recovering // tikv will get this state during BR EBS restore. // we write this info into etcd for simplicity, the key only stays inside etcd temporary From d2824bf11677debb88ed8c50a1a55b0024787c19 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 14:05:27 +0800 Subject: [PATCH 5/7] server: use PD mode terminology for metadata cleanup Signed-off-by: Ryan Leung --- server/api/admin.go | 4 ++-- server/microservice_cleanup_test.go | 2 +- server/server.go | 12 ++++++------ 3 files changed, 9 insertions(+), 9 deletions(-) diff --git a/server/api/admin.go b/server/api/admin.go index 81fd93f1c2..16e13e35d3 100644 --- a/server/api/admin.go +++ b/server/api/admin.go @@ -59,10 +59,10 @@ func newAdminHandler(svr *server.Server, rd *render.Render) *adminHandler { } // CleanupMicroserviceMetadata cleans up stale microservice runtime metadata in -// normal PD mode. +// PD mode. // // @Tags admin -// @Summary Clean up stale microservice runtime metadata in normal PD mode. +// @Summary Clean up stale microservice runtime metadata in PD mode. // @Produce json // @Success 200 {object} CleanupMicroserviceMetadataResponse "Whether any metadata was changed." // @Failure 409 {string} string "The cleanup is rejected in the current state." diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go index 5b4c37cdbb..bf333f2bfa 100644 --- a/server/microservice_cleanup_test.go +++ b/server/microservice_cleanup_test.go @@ -172,7 +172,7 @@ func TestCleanupMicroserviceMetadataRejectsUnsafeState(t *testing.T) { } } -func TestCleanupMicroserviceMetadataRequiresNormalPDLeader(t *testing.T) { +func TestCleanupMicroserviceMetadataRequiresPDModeLeader(t *testing.T) { t.Run("microservice-mode", func(t *testing.T) { svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13005) svr.isKeyspaceGroupEnabled = true diff --git a/server/server.go b/server/server.go index 303a94db51..e2b5bca553 100644 --- a/server/server.go +++ b/server/server.go @@ -2254,7 +2254,7 @@ var ( // keyspace-group state cannot be safely cleaned up. ErrMicroserviceMetadataCleanupRejected = errors.New("microservice metadata cleanup rejected") // ErrMicroserviceMetadataCleanupUnavailable indicates that the request could - // not be completed under the current normal-PD leadership term. + // not be completed under the current leadership term in PD mode. ErrMicroserviceMetadataCleanupUnavailable = errors.New("microservice metadata cleanup unavailable") ) @@ -2273,8 +2273,8 @@ func unavailableMicroserviceMetadataCleanup(format string, args ...any) error { } // CleanupMicroserviceMetadata clears the persisted Members field of the default -// TSO keyspace group in normal PD mode and reports whether it changed. A nil -// error is fenced by one exact normal-PD leadership term and represents a +// TSO keyspace group in PD mode and reports whether it changed. A nil error is +// fenced by one exact leadership term in PD mode and represents a // linearizable check that the default group existed, had no transition in // progress, had no non-default sibling, and had no persisted member at that // point. The caller must ensure keyspace assignments have already been merged @@ -2371,12 +2371,12 @@ func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataCleanupTerm, error) { if s.client == nil || s.member == nil || !s.member.IsServing() { return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leader is not serving") + "leader in PD mode is not serving") } leadership := s.member.GetLeadership() if leadership == nil || leadership.GetLease() == nil { return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leadership is not initialized") + "leadership in PD mode is not initialized") } term := microserviceMetadataCleanupTerm{ leaderKey: leadership.GetLeaderKey(), @@ -2385,7 +2385,7 @@ func (s *Server) captureMicroserviceMetadataCleanupTerm() (microserviceMetadataC } if term.leaderKey == "" || term.leaderValue == "" || term.leaseID == 0 { return microserviceMetadataCleanupTerm{}, unavailableMicroserviceMetadataCleanup( - "normal PD leadership term is incomplete") + "leadership term in PD mode is incomplete") } return term, nil } From e302f559400684b749e8c3bb72b0d6a75444b7a3 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 16:19:51 +0800 Subject: [PATCH 6/7] server: preserve unknown metadata during cleanup Rewrite only the members field in the persisted keyspace group JSON so additive fields written by newer versions survive cleanup. Reject ambiguous or malformed group objects before returning a successful cleanup certificate. Signed-off-by: Ryan Leung --- server/microservice_cleanup_test.go | 67 +++++++++++++++++++++++++++++ server/server.go | 63 +++++++++++++++++++++++++-- 2 files changed, 127 insertions(+), 3 deletions(-) diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go index bf333f2bfa..3439ac8c87 100644 --- a/server/microservice_cleanup_test.go +++ b/server/microservice_cleanup_test.go @@ -15,7 +15,9 @@ package server import ( + "bytes" "context" + "encoding/json" "sync" "testing" "time" @@ -94,6 +96,71 @@ func TestCleanupMicroserviceMetadataPreservesDefaultGroup(t *testing.T) { require.Equal(t, timestampValue, string(resp.Kvs[0].Value)) } +func TestCleanupMicroserviceMetadataPreservesUnknownGroupFields(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13015) + ctx := context.Background() + const ( + futureField = `{"large-id":9007199254740993,"nested":{"enabled":true}}` + groupJSON = `{"id":0,"user-kind":"basic","members":[{"address":"http://127.0.0.1:3379","priority":0}],"keyspaces":[1],"future-field":` + futureField + `}` + ) + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) + _, err := svr.client.Put(ctx, groupKey, groupJSON) + require.NoError(t, err) + + changed, err := svr.CleanupMicroserviceMetadata(ctx) + require.NoError(t, err) + require.True(t, changed) + + resp, err := etcdutil.EtcdKVGet(svr.client, groupKey) + require.NoError(t, err) + require.Len(t, resp.Kvs, 1) + persistedGroup := make(map[string]json.RawMessage) + require.NoError(t, json.Unmarshal(resp.Kvs[0].Value, &persistedGroup)) + require.Equal(t, "null", string(persistedGroup["members"])) + require.True(t, bytes.Equal(json.RawMessage(futureField), persistedGroup["future-field"])) + + value := string(resp.Kvs[0].Value) + modRevision := resp.Kvs[0].ModRevision + changed, err = svr.CleanupMicroserviceMetadata(ctx) + require.NoError(t, err) + require.False(t, changed) + resp, err = etcdutil.EtcdKVGet(svr.client, groupKey) + require.NoError(t, err) + require.Len(t, resp.Kvs, 1) + require.Equal(t, value, string(resp.Kvs[0].Value)) + require.Equal(t, modRevision, resp.Kvs[0].ModRevision) +} + +func TestCleanupMicroserviceMetadataRejectsInvalidGroupJSONObjects(t *testing.T) { + testCases := []struct { + name string + groupJSON string + }{ + {name: "null", groupJSON: "null"}, + {name: "missing-id", groupJSON: `{"members":[]}`}, + {name: "non-canonical-members", groupJSON: `{"id":0,"Members":[{"address":"http://127.0.0.1:3379","priority":0}]}`}, + {name: "duplicate-members", groupJSON: `{"id":0,"members":[{"address":"http://127.0.0.1:3379","priority":0}],"members":[]}`}, + {name: "duplicate-transition", groupJSON: `{"id":0,"split-state":{"split-source":0},"split-state":null,"members":[]}`}, + } + for i, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, uint64(13016+i)) + ctx := context.Background() + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) + _, err := svr.client.Put(ctx, groupKey, testCase.groupJSON) + require.NoError(t, err) + + changed, err := svr.CleanupMicroserviceMetadata(ctx) + require.False(t, changed) + require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + resp, err := etcdutil.EtcdKVGet(svr.client, groupKey) + require.NoError(t, err) + require.Len(t, resp.Kvs, 1) + require.Equal(t, testCase.groupJSON, string(resp.Kvs[0].Value)) + }) + } +} + func TestCleanupMicroserviceMetadataRejectsUnsafeState(t *testing.T) { t.Run("missing-default-group", func(t *testing.T) { svr, store := newMicroserviceMetadataCleanupTestServer(t, 13014) diff --git a/server/server.go b/server/server.go index e2b5bca553..38316d83d8 100644 --- a/server/server.go +++ b/server/server.go @@ -2310,11 +2310,66 @@ func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) return false, rejectMicroserviceMetadataCleanup("found a non-default TSO keyspace group") } + persistedGroup := make(map[string]json.RawMessage) + if err := json.Unmarshal(resp.Kvs[0].Value, &persistedGroup); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + if persistedGroup == nil { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group metadata is not a JSON object") + } + decoder := json.NewDecoder(bytes.NewReader(resp.Kvs[0].Value)) + if _, err := decoder.Token(); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + seenFields := make(map[string]struct{}, len(persistedGroup)) + for decoder.More() { + token, err := decoder.Token() + if err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + field, ok := token.(string) + if !ok { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group metadata has a non-string field") + } + if _, ok := seenFields[field]; ok { + return false, rejectMicroserviceMetadataCleanup( + "default TSO keyspace group metadata contains duplicate field %q", field) + } + seenFields[field] = struct{}{} + var value json.RawMessage + if err := decoder.Decode(&value); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + } + if _, err := decoder.Token(); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + groupIDValue, ok := persistedGroup["id"] + if !ok { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group metadata has no ID") + } + var persistedGroupID *uint32 + if err := json.Unmarshal(groupIDValue, &persistedGroupID); err != nil { + return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() + } + if persistedGroupID == nil { + return false, rejectMicroserviceMetadataCleanup("default TSO keyspace group metadata has no ID") + } + canonicalFields := [...]string{"id", "user-kind", "split-state", "merge-state", "members", "keyspaces"} + for key := range persistedGroup { + for _, canonicalField := range canonicalFields { + if key != canonicalField && strings.EqualFold(key, canonicalField) { + return false, rejectMicroserviceMetadataCleanup( + "default TSO keyspace group metadata contains non-canonical field %q", key) + } + } + } + group := &endpoint.KeyspaceGroup{} if err := json.Unmarshal(resp.Kvs[0].Value, group); err != nil { return false, errs.ErrJSONUnmarshal.Wrap(err).GenWithStackByCause() } - if group.ID != constant.DefaultKeyspaceGroupID { + if *persistedGroupID != constant.DefaultKeyspaceGroupID || group.ID != constant.DefaultKeyspaceGroupID { return false, rejectMicroserviceMetadataCleanup( "found TSO keyspace group %d at the default group path", group.ID) } @@ -2344,8 +2399,10 @@ func (s *Server) CleanupMicroserviceMetadata(ctx context.Context) (bool, error) } operation := clientv3.OpGet(groupKey) if changed { - group.Members = nil - value, err := json.Marshal(group) + // Rewrite only the members field so metadata written by a newer version + // is not lost if it contains fields this binary does not recognize. + persistedGroup["members"] = json.RawMessage("null") + value, err := json.Marshal(persistedGroup) if err != nil { return false, errs.ErrJSONMarshal.Wrap(err).GenWithStackByCause() } From 9eafa3589fc326a34281b1ed55fe8f10fef0dee0 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 13 Aug 2026 16:58:59 +0800 Subject: [PATCH 7/7] tests: reduce metadata cleanup test duplication Signed-off-by: Ryan Leung --- server/microservice_cleanup_test.go | 47 ++++++++--------------- tests/integrations/mcs/tso/server_test.go | 20 +--------- tests/server/api/admin_test.go | 7 +--- 3 files changed, 20 insertions(+), 54 deletions(-) diff --git a/server/microservice_cleanup_test.go b/server/microservice_cleanup_test.go index 3439ac8c87..fa2aca7a7c 100644 --- a/server/microservice_cleanup_test.go +++ b/server/microservice_cleanup_test.go @@ -78,10 +78,6 @@ func TestCleanupMicroserviceMetadataPreservesDefaultGroup(t *testing.T) { require.Nil(t, group.SplitState) require.Nil(t, group.MergeState) - changed, err = svr.CleanupMicroserviceMetadata(ctx) - require.NoError(t, err) - require.False(t, changed) - require.NoError(t, store.RunInTxn(ctx, func(txn kv.Txn) error { meta, err := store.LoadKeyspaceMeta(txn, 1) require.NoError(t, err) @@ -132,6 +128,9 @@ func TestCleanupMicroserviceMetadataPreservesUnknownGroupFields(t *testing.T) { } func TestCleanupMicroserviceMetadataRejectsInvalidGroupJSONObjects(t *testing.T) { + svr, _ := newMicroserviceMetadataCleanupTestServer(t, 13016) + ctx := context.Background() + groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) testCases := []struct { name string groupJSON string @@ -142,11 +141,8 @@ func TestCleanupMicroserviceMetadataRejectsInvalidGroupJSONObjects(t *testing.T) {name: "duplicate-members", groupJSON: `{"id":0,"members":[{"address":"http://127.0.0.1:3379","priority":0}],"members":[]}`}, {name: "duplicate-transition", groupJSON: `{"id":0,"split-state":{"split-source":0},"split-state":null,"members":[]}`}, } - for i, testCase := range testCases { + for _, testCase := range testCases { t.Run(testCase.name, func(t *testing.T) { - svr, _ := newMicroserviceMetadataCleanupTestServer(t, uint64(13016+i)) - ctx := context.Background() - groupKey := keypath.KeyspaceGroupIDPath(constant.DefaultKeyspaceGroupID) _, err := svr.client.Put(ctx, groupKey, testCase.groupJSON) require.NoError(t, err) @@ -284,9 +280,7 @@ func TestCleanupMicroserviceMetadataIsFencedByExactLeadershipTerm(t *testing.T) blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + requireMicroserviceMetadataCleanupCASConflict(t, result) require.NotEmpty(t, loadMicroserviceMetadataCleanupTestGroup( t, store, constant.DefaultKeyspaceGroupID).Members) } @@ -310,9 +304,7 @@ func TestCleanupMicroserviceMetadataPreservesConcurrentGroupUpdate(t *testing.T) blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + requireMicroserviceMetadataCleanupCASConflict(t, result) group := loadMicroserviceMetadataCleanupTestGroup(t, store, constant.DefaultKeyspaceGroupID) require.Equal(t, newMembers, group.Members) require.Equal(t, []uint32{1, 2}, group.Keyspaces) @@ -332,9 +324,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + requireMicroserviceMetadataCleanupCASConflict(t, result) }) t.Run("empty-members-group-update", func(t *testing.T) { @@ -352,9 +342,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + requireMicroserviceMetadataCleanupCASConflict(t, result) require.Equal(t, []uint32{1, 2}, loadMicroserviceMetadataCleanupTestGroup( t, store, constant.DefaultKeyspaceGroupID).Keyspaces) }) @@ -374,13 +362,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") - - changed, err := svr.CleanupMicroserviceMetadata(context.Background()) - require.False(t, changed) - require.ErrorIs(t, err, ErrMicroserviceMetadataCleanupRejected) + requireMicroserviceMetadataCleanupCASConflict(t, result) }) t.Run("non-default-created", func(t *testing.T) { @@ -399,9 +381,7 @@ func TestCleanupMicroserviceMetadataFencesNoOp(t *testing.T) { blocker.releaseCleanup() result := waitMicroserviceMetadataCleanupResult(t, resultCh) - require.False(t, result.changed) - require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) - require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") + requireMicroserviceMetadataCleanupCASConflict(t, result) require.NotNil(t, loadMicroserviceMetadataCleanupTestGroup(t, store, 1)) }) } @@ -413,6 +393,13 @@ type microserviceMetadataCleanupResult struct { err error } +func requireMicroserviceMetadataCleanupCASConflict(t *testing.T, result microserviceMetadataCleanupResult) { + t.Helper() + require.False(t, result.changed) + require.ErrorIs(t, result.err, ErrMicroserviceMetadataCleanupUnavailable) + require.ErrorContains(t, result.err, "leadership or keyspace-group metadata changed during cleanup") +} + type microserviceMetadataCleanupCommitBlocker struct { name string reached chan struct{} diff --git a/tests/integrations/mcs/tso/server_test.go b/tests/integrations/mcs/tso/server_test.go index 8de7fb7b99..e7914f704e 100644 --- a/tests/integrations/mcs/tso/server_test.go +++ b/tests/integrations/mcs/tso/server_test.go @@ -809,7 +809,7 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { oldTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) re.NoError(err) defer oldTSOCluster.Destroy() - oldTSOPrimary := waitForTSOServiceReady(re, oldTSOCluster) + oldTSOPrimary := oldTSOCluster.WaitForDefaultPrimaryServing(re) oldMembers := oldTSOCluster.GetKeyspaceGroupMember() waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) @@ -829,13 +829,11 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { oldTSOCluster.Destroy() assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) pdServer = restartPDWithServices(ctx, re, tc, pdServer, nil) waitForTSOMonotonic(ctx, re, defaultClient, &globalLastTS) waitForTSOMonotonic(ctx, re, userClient, &globalLastTS) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, oldMembers) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) re.True(cleanupMicroserviceMetadataViaHTTP(re, pdServer)) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) @@ -845,12 +843,11 @@ func TestCleanupMicroserviceMetadataForModeSwitch(t *testing.T) { pdServer = restartPDWithServices(ctx, re, tc, pdServer, []string{mcs.PDServiceName}) assertDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, nil) - assertUserKeyspaceInDefaultGroup(re, pdServer, userKeyspaceID) newTSOCluster, err := tests.NewTestTSOCluster(ctx, 1, pdServer.GetAddr()) re.NoError(err) defer newTSOCluster.Destroy() - waitForTSOServiceReady(re, newTSOCluster) + newTSOCluster.WaitForDefaultPrimaryServing(re) newMembers := newTSOCluster.GetKeyspaceGroupMember() waitForDefaultKeyspaceGroup(re, pdServer, userKeyspaceID, newMembers) waitForTSOMonotonic(ctx, re, defaultClient, &globalLastTS) @@ -957,19 +954,6 @@ func assertUserKeyspaceInDefaultGroup( re.Equal("0", meta.GetConfig()[keyspace.TSOKeyspaceGroupIDKey]) } -func waitForTSOServiceReady(re *require.Assertions, tsoCluster *tests.TestTSOCluster) *tso.Server { - 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)) - return primary -} - func waitForTSOMonotonic(ctx context.Context, re *require.Assertions, client pd.Client, globalLastTS *uint64) { for range 10 { var physical, logical int64 diff --git a/tests/server/api/admin_test.go b/tests/server/api/admin_test.go index fd1c225cbd..6aaff8fcdc 100644 --- a/tests/server/api/admin_test.go +++ b/tests/server/api/admin_test.go @@ -409,20 +409,15 @@ func (suite *adminTestSuite) checkCleanupMicroserviceMetadata(cluster *tests.Tes })) re.NoError(testutil.CheckGetJSON(tests.TestDialClient, url, nil, testutil.Status(re, http.StatusMethodNotAllowed))) - storedGroups, err := storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) - re.NoError(err) - re.Len(storedGroups, 1) - re.Equal(group.Members, storedGroups[0].Members) response := &api.CleanupMicroserviceMetadataResponse{} re.NoError(testutil.CheckPostJSON(tests.TestDialClient, url, nil, testutil.StatusOK(re), testutil.ExtractJSON(re, response))) re.True(response.Changed) - storedGroups, err = storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) + storedGroups, err := storage.LoadKeyspaceGroups(keyspaceconstant.DefaultKeyspaceGroupID, 1) re.NoError(err) re.Len(storedGroups, 1) re.Empty(storedGroups[0].Members) - re.Equal(group.Keyspaces, storedGroups[0].Keyspaces) response.Changed = true re.NoError(testutil.CheckPostJSON(tests.TestDialClient, url, nil,