Skip to content
Draft
2 changes: 1 addition & 1 deletion api/middleware/authenticate_middleware.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ func verify(ctx *gin.Context, etcdCli etcd.Client) error {

// verifyTiDBUser verify whether the username and password are valid in TiDB. It does the validation via
// the successfully build of a connection with upstream TiDB with the username and password.
tidbs, err := upstream.FetchTiDBTopology(ctx, etcdCli, keyspaceMeta.Id)
tidbs, err := upstream.FetchTiDBTopology(ctx, etcdCli, keyspaceMeta.GetId())
if err != nil {
return errors.Trace(err)
}
Expand Down
12 changes: 6 additions & 6 deletions api/middleware/middleware_nextgen_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,18 +75,18 @@ func TestKeyspaceCheckerMiddleware(t *testing.T) {
keyspace: "success",
init: func(t *testing.T, mock *keyspace.MockManager) {
mock.EXPECT().LoadKeyspace(gomock.Any(), "success").Return(&keyspacepb.KeyspaceMeta{
Id: 1,
Name: "kespace1",
State: keyspacepb.KeyspaceState_ENABLED,
Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: 1},
Name: "kespace1",
State: keyspacepb.KeyspaceState_ENABLED,
}, nil)
},
expectedStatus: http.StatusOK,
expectedAbort: false,
expectedBodyContains: "",
expectedMeta: &keyspacepb.KeyspaceMeta{
Id: 1,
Name: "kespace1",
State: keyspacepb.KeyspaceState_ENABLED,
Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: 1},
Name: "kespace1",
State: keyspacepb.KeyspaceState_ENABLED,
},
},
}
Expand Down
12 changes: 6 additions & 6 deletions api/v2/changefeed.go
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ func (h *OpenAPIV2) CreateChangefeed(c *gin.Context) {
ctx,
h.server.GetPdClient(),
createGcServiceID,
keyspaceMeta.Id,
keyspaceMeta.GetId(),
changefeedID,
ensureTTL, cfg.StartTs); err != nil {
if !errors.ErrStartTsBeforeGC.Equal(err) {
Expand All @@ -240,7 +240,7 @@ func (h *OpenAPIV2) CreateChangefeed(c *gin.Context) {
undoErr := gc.UndoEnsureChangefeedStartTsSafety(
ctx,
pdClient,
keyspaceMeta.Id,
keyspaceMeta.GetId(),
createGcServiceID,
changefeedID,
)
Expand Down Expand Up @@ -268,7 +268,7 @@ func (h *OpenAPIV2) CreateChangefeed(c *gin.Context) {
// We create a new context here.
schemaCxt := context.Background()
if err = schemaStore.RegisterKeyspace(schemaCxt, common.KeyspaceMeta{
ID: keyspaceMeta.Id,
ID: keyspaceMeta.GetId(),
Name: keyspaceMeta.Name,
}); err != nil {
_ = c.Error(err)
Expand Down Expand Up @@ -304,7 +304,7 @@ func (h *OpenAPIV2) CreateChangefeed(c *gin.Context) {
Config: replicaCfg,
State: config.StateNormal,
CreatorVersion: version.ReleaseVersion,
KeyspaceID: keyspaceMeta.Id,
KeyspaceID: keyspaceMeta.GetId(),
}

// verify sinkURI
Expand Down Expand Up @@ -826,7 +826,7 @@ func (h *OpenAPIV2) ResumeChangefeed(c *gin.Context) {
ctx,
h.server.GetPdClient(),
resumeGcServiceID,
keyspaceMeta.Id,
keyspaceMeta.GetId(),
cfInfo.ChangefeedID,
newCheckpointTs); err != nil {
_ = c.Error(err)
Expand All @@ -840,7 +840,7 @@ func (h *OpenAPIV2) ResumeChangefeed(c *gin.Context) {
undoErr := gc.UndoEnsureChangefeedStartTsSafety(
ctx,
h.server.GetPdClient(),
keyspaceMeta.Id,
keyspaceMeta.GetId(),
resumeGcServiceID,
cfInfo.ChangefeedID,
)
Expand Down
4 changes: 2 additions & 2 deletions api/v2/changefeed_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,8 +70,8 @@ func TestResumeChangefeedRejectsNormalBeforeGC(t *testing.T) {
c.Request = httptest.NewRequest(http.MethodPost, "/api/v2/changefeeds/test/resume?keyspace=default", nil)
c.Params = gin.Params{{Key: api.APIOpVarChangefeedID, Value: "test"}}
c.Set("ctx-keyspace", &keyspacepb.KeyspaceMeta{
Id: common.DefaultKeyspaceID,
State: keyspacepb.KeyspaceState_ENABLED,
Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: common.DefaultKeyspaceID},
State: keyspacepb.KeyspaceState_ENABLED,
})

h.ResumeChangefeed(c)
Expand Down
6 changes: 3 additions & 3 deletions api/v2/unsafe.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,10 +59,10 @@ func (h *OpenAPIV2) ResolveLock(c *gin.Context) {
keyspaceMeta := middleware.GetKeyspaceFromContext(c)

txnResolver := txnutil.NewLockerResolver()
if err := txnResolver.Resolve(schemaCxt, keyspaceMeta.Id, resolveLockReq.RegionID, resolveLockReq.Ts); err != nil {
if err := txnResolver.Resolve(schemaCxt, keyspaceMeta.GetId(), resolveLockReq.RegionID, resolveLockReq.Ts); err != nil {
log.Error(
"resolve lock failed",
zap.Uint32("keyspaceID", keyspaceMeta.Id),
zap.Uint32("keyspaceID", keyspaceMeta.GetId()),
zap.Uint64("regionID", resolveLockReq.RegionID),
zap.Uint64("resolveLockTs", resolveLockReq.Ts),
zap.Error(err),
Expand All @@ -87,7 +87,7 @@ func (h *OpenAPIV2) DeleteServiceGcSafePoint(c *gin.Context) {
err := gc.UnifyDeleteGcSafepoint(
c,
pdClient,
keyspaceMeta.Id,
keyspaceMeta.GetId(),
h.server.GetEtcdClient().GetGCServiceID(),
)
if err != nil {
Expand Down
1 change: 1 addition & 0 deletions cmd/cdc/cli/cli_unsafe.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ func newCmdUnsafe(f factory.Factory) *cobra.Command {
command.AddCommand(newCmdShowMetadata(f))
command.AddCommand(newCmdDeleteServiceGcSafepoint(f, commonOptions))
command.AddCommand(newCmdResolveLock(f))
command.AddCommand(newCmdVerifyGCSafepoint(f))

return command
}
4 changes: 2 additions & 2 deletions cmd/cdc/cli/cli_unsafe_reset.go
Original file line number Diff line number Diff line change
Expand Up @@ -139,9 +139,9 @@ func removeKeyspaceGCBarrier(ctx context.Context, pdCli pd.Client, serviceID str
log.Warn("load keyspace error", zap.String("keyspace", keyspace), zap.Error(err))
continue
}
err = gc.UnifyDeleteGcSafepoint(ctx, pdCli, keyspaceMeta.Id, serviceID)
err = gc.UnifyDeleteGcSafepoint(ctx, pdCli, keyspaceMeta.GetId(), serviceID)
if err != nil {
log.Warn("DeleteGcSafepoint error", zap.Uint32("keyspaceID", keyspaceMeta.Id), zap.String("keyspace", keyspace), zap.String("serviceID", serviceID), zap.Error(err))
log.Warn("DeleteGcSafepoint error", zap.Uint32("keyspaceID", keyspaceMeta.GetId()), zap.String("keyspace", keyspace), zap.String("serviceID", serviceID), zap.Error(err))
continue
}
}
Expand Down
159 changes: 159 additions & 0 deletions cmd/cdc/cli/cli_unsafe_verify_gc_safepoint.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
// Copyright 2026 PingCAP, Inc.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// See the License for the specific language governing permissions and
// limitations under the License.

package cli

import (
"context"
"net/url"

"github.com/pingcap/kvproto/pkg/keyspacepb"
"github.com/pingcap/ticdc/cmd/cdc/factory"
"github.com/pingcap/ticdc/cmd/util"
"github.com/pingcap/ticdc/pkg/common"
"github.com/pingcap/ticdc/pkg/errors"
"github.com/pingcap/ticdc/pkg/security"
"github.com/pingcap/ticdc/pkg/upstream"
tidbkv "github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/meta"
"github.com/pingcap/tidb/pkg/store/driver"
"github.com/spf13/cobra"
tikvconfig "github.com/tikv/client-go/v2/config"
pdgc "github.com/tikv/pd/client/clients/gc"
)

type verifyGCSafepointPDClient interface {
LoadKeyspace(ctx context.Context, keyspace string) (*keyspacepb.KeyspaceMeta, error)
GetGCStatesClient(keyspaceID uint32) pdgc.GCStatesClient
Close()
}

type verifyGCSafepointOptions struct {
keyspace string
legacySafepoint bool
pdClient verifyGCSafepointPDClient
listDatabases func(ctx context.Context, keyspace string, snapshotTS uint64) (int, error)
}

func (o *verifyGCSafepointOptions) addFlags(cmd *cobra.Command) {
cmd.Flags().StringVarP(&o.keyspace, "keyspace", "k", "", "Keyspace to verify")
cmd.Flags().BoolVar(&o.legacySafepoint, "legacy-safepoint", false, "Read the keyspace-v2 minimum service safepoint")
_ = cmd.MarkFlagRequired("keyspace")
}

func (o *verifyGCSafepointOptions) complete(f factory.Factory) error {
pdClient, err := f.PdClient()
if err != nil {
return err
}
o.pdClient = pdClient
o.listDatabases = func(_ context.Context, keyspace string, snapshotTS uint64) (int, error) {
return listDatabasesAtSnapshot(f.GetPdAddr(), f.GetCredential(), keyspace, snapshotTS)
}
return nil
}

func (o *verifyGCSafepointOptions) run(cmd *cobra.Command) error {
defer o.pdClient.Close()
ctx := cmd.Context()
keyspaceMeta, err := o.pdClient.LoadKeyspace(ctx, o.keyspace)
if err != nil {
return errors.WrapError(errors.ErrLoadKeyspaceFailed, err)
}

var snapshotTS uint64
if o.legacySafepoint {
legacyClient, ok := o.pdClient.(pdgc.LegacyClientV2)
if !ok {
return errors.ErrGetServiceSafepointFailed.GenWithStackByArgs("PD client does not support LegacyClientV2")
}
snapshotTS, err = legacyClient.GetMinServiceSafePointV2(ctx, keyspaceMeta.GetId())
if err != nil {
return errors.WrapError(errors.ErrGetServiceSafepointFailed, err)
}
} else {
gcState, err := o.pdClient.GetGCStatesClient(keyspaceMeta.GetId()).GetGCState(ctx)
if err != nil {
return errors.Trace(err)
}
// Schema store uses the transaction safepoint as its initial metadata snapshot.
snapshotTS = gcState.TxnSafePoint
}
if snapshotTS == 0 {
return errors.ErrGetServiceSafepointFailed.GenWithStackByArgs(
"safepoint is zero, keyspace GC state may be uninitialized")
}
cmd.Printf("WARNING: this verification holds no service safepoint; GC may advance past snapshot %d before the read completes.\n", snapshotTS)

databaseCount, err := o.listDatabases(ctx, o.keyspace, snapshotTS)
if err != nil {
return err
}
cmd.Printf("GC safepoint and ListDatabases verified, snapshot-ts: %d, databases: %d\n", snapshotTS, databaseCount)
return nil
}

func listDatabasesAtSnapshot(
pdAddr string,
credential *security.Credential,
keyspace string,
snapshotTS uint64,
) (int, error) {
pdURLs, err := common.NewURLsValue(pdAddr)
if err != nil {
return 0, errors.WrapError(errors.ErrNewStore, err)
}
tiURL := &url.URL{Scheme: "tikv", Host: pdURLs.HostString()}
query := tiURL.Query()
query.Set("disableGC", "true")
query.Set("keyspaceName", keyspace)
tiURL.RawQuery = query.Encode()

securityConfig := tikvconfig.Security{
ClusterSSLCA: credential.CAPath,
ClusterSSLCert: credential.CertPath,
ClusterSSLKey: credential.KeyPath,
ClusterVerifyCN: credential.CertAllowedCN,
}
tiStore, err := (&driver.TiKVDriver{}).OpenWithOptions(
tiURL.String(), driver.WithSecurity(securityConfig),
)
if err != nil {
return 0, errors.WrapError(errors.ErrNewStore, err)
}
defer func() { _ = tiStore.Close() }()
if err := upstream.DisablePDRouterClient(tiStore); err != nil {
return 0, errors.WrapError(errors.ErrNewStore, err)
}

databases, err := meta.NewReader(tiStore.GetSnapshot(tidbkv.NewVersion(snapshotTS))).ListDatabases()
if err != nil {
return 0, errors.WrapError(errors.ErrMetaListDatabases, err)
}
return len(databases), nil
}

func newCmdVerifyGCSafepoint(f factory.Factory) *cobra.Command {
o := &verifyGCSafepointOptions{}
command := &cobra.Command{
Use: "verify-gc-safepoint",
Short: "Verify a keyspace GC safepoint by listing databases at its current snapshot",
Args: cobra.NoArgs,
Run: func(cmd *cobra.Command, _ []string) {
util.CheckErr(o.complete(f))
util.CheckErr(o.run(cmd))
},
}
o.addFlags(command)
return command
}
Loading
Loading