diff --git a/examples/gcworker/go.mod b/examples/gcworker/go.mod index 26d15ac7c6..b4be0c1094 100644 --- a/examples/gcworker/go.mod +++ b/examples/gcworker/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/rawkv/go.mod b/examples/rawkv/go.mod index f7d842531f..a18ccbf133 100644 --- a/examples/rawkv/go.mod +++ b/examples/rawkv/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/1pc_txn/go.mod b/examples/txnkv/1pc_txn/go.mod index 4655c6bac2..fb7c71f85d 100644 --- a/examples/txnkv/1pc_txn/go.mod +++ b/examples/txnkv/1pc_txn/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/async_commit/go.mod b/examples/txnkv/async_commit/go.mod index 7da0131aea..e95b288cf1 100644 --- a/examples/txnkv/async_commit/go.mod +++ b/examples/txnkv/async_commit/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/delete_range/go.mod b/examples/txnkv/delete_range/go.mod index 46cc6dc311..1b38f46343 100644 --- a/examples/txnkv/delete_range/go.mod +++ b/examples/txnkv/delete_range/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/go.mod b/examples/txnkv/go.mod index 61df27ff81..88600ffe72 100644 --- a/examples/txnkv/go.mod +++ b/examples/txnkv/go.mod @@ -3,7 +3,7 @@ module txnkv go 1.25.10 require ( - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 github.com/tikv/client-go/v2 v2.0.0 ) @@ -33,7 +33,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/pessimistic_txn/go.mod b/examples/txnkv/pessimistic_txn/go.mod index 31d796f9a4..ba1cfabd64 100644 --- a/examples/txnkv/pessimistic_txn/go.mod +++ b/examples/txnkv/pessimistic_txn/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/examples/txnkv/unsafedestoryrange/go.mod b/examples/txnkv/unsafedestoryrange/go.mod index 1538fe6892..02b223b259 100644 --- a/examples/txnkv/unsafedestoryrange/go.mod +++ b/examples/txnkv/unsafedestoryrange/go.mod @@ -22,7 +22,7 @@ require ( github.com/opentracing/opentracing-go v1.2.0 // indirect github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 // indirect github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 // indirect - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 // indirect + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 // indirect github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/prometheus/client_golang v1.20.5 // indirect @@ -31,7 +31,7 @@ require ( github.com/prometheus/procfs v0.15.1 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a // indirect - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 // indirect + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f // indirect github.com/twmb/murmur3 v1.1.3 // indirect go.etcd.io/etcd/api/v3 v3.5.10 // indirect go.etcd.io/etcd/client/pkg/v3 v3.5.10 // indirect diff --git a/go.mod b/go.mod index 4578fd7b0b..f87b5c393f 100644 --- a/go.mod +++ b/go.mod @@ -15,14 +15,14 @@ require ( github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.20.5 github.com/prometheus/client_model v0.6.1 github.com/stretchr/testify v1.9.0 github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f github.com/twmb/murmur3 v1.1.3 go.etcd.io/etcd/api/v3 v3.5.10 go.etcd.io/etcd/client/v3 v3.5.10 diff --git a/go.sum b/go.sum index 3bde9f254d..78c5053e4d 100644 --- a/go.sum +++ b/go.sum @@ -81,8 +81,8 @@ github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 h1:tdMsjOqUR7YXH github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86/go.mod h1:exzhVYca3WRtd6gclGNErRWb1qEgff3LYta0LvRmON4= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 h1:surzm05a8C9dN8dIUmo4Be2+pMRb6f55i+UIYrluu2E= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989/go.mod h1:O17XtbryoCJhkKGbT62+L2OlrniwqiGLSqrmdHCMzZw= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 h1:qW0gHsqY3X3qmyiEfESqpYSF3Vu7agUiAyBnPWQXtm8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= +github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 h1:JBrbgynAz4RkeXUvqXAOhySZK5FYFCeHdKWUrj+Z2Vw= +github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 h1:HR/ylkkLmGdSSDaD8IDP+SZrdhV1Kibl9KrHxJ9eciw= github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= @@ -115,8 +115,8 @@ github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsT github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a h1:J/YdBZ46WKpXsxsW93SG+q0F8KI+yFrcIDT4c/RNoc4= github.com/tiancaiamao/gp v0.0.0-20221230034425-4025bc8a4d4a/go.mod h1:h4xBhSNtOeEosLJ4P7JyKXX7Cabg7AVkWCK5gV2vOrM= -github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 h1:OoBvgoeWmdNEXtS+eOlhysz/OvhA4GS0OdPVhTXteGA= -github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3/go.mod h1:3/Bu91CJONgkDA+Y0v/cnbROSJnu5tQ09vv7JGybUBA= +github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f h1:IYHVxTMV8oXlL4ZaoGDvREUZlszfyA0w3fUK9g2wSlM= +github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f/go.mod h1:Xhd/CPuOv+tgP+5LKhWS23tklCXszF0EgKo4R6OZGj0= github.com/twmb/murmur3 v1.1.3 h1:D83U0XYKcHRYwYIpBKf3Pks91Z0Byda/9SJ8B6EMRcA= github.com/twmb/murmur3 v1.1.3/go.mod h1:Qq/R7NUyOfr65zD+6Q5IHKsJLwP7exErjN6lyyq3OSQ= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= diff --git a/integration_tests/gc_test.go b/integration_tests/gc_test.go index fdc9c0df9d..b4cf042ab0 100644 --- a/integration_tests/gc_test.go +++ b/integration_tests/gc_test.go @@ -179,7 +179,7 @@ func (s *testGCWithTiKVSuite) dropKeyspace(keyspaceMeta *keyspacepb.KeyspaceMeta re := s.Require() // Nil might be used to represent the null keyspace. if keyspaceMeta != nil { - _, err := s.globalPDCli.UpdateKeyspaceState(context.Background(), keyspaceMeta.Id, keyspacepb.KeyspaceState_ARCHIVED) + _, err := s.globalPDCli.UpdateKeyspaceState(context.Background(), keyspaceMeta.GetId(), keyspacepb.KeyspaceState_ARCHIVED) re.NoError(err) } } diff --git a/integration_tests/go.mod b/integration_tests/go.mod index 02c7cf6571..cc0ae245bd 100644 --- a/integration_tests/go.mod +++ b/integration_tests/go.mod @@ -7,7 +7,7 @@ require ( github.com/ninedraft/israce v0.0.3 github.com/pingcap/errors v0.11.5-0.20260508054701-306e305bcf41 github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 + github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 github.com/pingcap/tidb v1.1.0-beta.0.20260715060322-10292a4f8697 github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.23.0 @@ -15,7 +15,7 @@ require ( github.com/stretchr/testify v1.11.1 github.com/tidwall/gjson v1.14.4 github.com/tikv/client-go/v2 v2.0.8-0.20260708122311-01bd8f99f4da - github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 + github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f go.uber.org/goleak v1.3.0 google.golang.org/grpc v1.79.3 ) diff --git a/integration_tests/go.sum b/integration_tests/go.sum index f0e172e567..72e34e3e18 100644 --- a/integration_tests/go.sum +++ b/integration_tests/go.sum @@ -1452,8 +1452,8 @@ github.com/pingcap/fn v1.0.0/go.mod h1:u9WZ1ZiOD1RpNhcI42RucFh/lBuzTu6rw88a+oF2Z github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 h1:surzm05a8C9dN8dIUmo4Be2+pMRb6f55i+UIYrluu2E= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989/go.mod h1:O17XtbryoCJhkKGbT62+L2OlrniwqiGLSqrmdHCMzZw= github.com/pingcap/kvproto v0.0.0-20241113043844-e1fa7ea8c302/go.mod h1:rXxWk2UnwfUhLXha1jxRWPADw9eMZGWEWCg92Tgmb/8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 h1:qW0gHsqY3X3qmyiEfESqpYSF3Vu7agUiAyBnPWQXtm8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= +github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472 h1:JBrbgynAz4RkeXUvqXAOhySZK5FYFCeHdKWUrj+Z2Vw= +github.com/pingcap/kvproto v0.0.0-20260724054804-059694ae4472/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/pingcap/log v0.0.0-20210625125904-98ed8e2eb1c7/go.mod h1:8AanEdAHATuRurdGxZXBz0At+9avep+ub7U1AGYLIMM= github.com/pingcap/log v1.1.0/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4= github.com/pingcap/log v1.1.1-0.20250917021125-19901e015dc9 h1:qG9BSvlWFEE5otQGamuWedx9LRm0nrHvsQRQiW8SxEs= @@ -1598,8 +1598,8 @@ github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JT github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= -github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 h1:OoBvgoeWmdNEXtS+eOlhysz/OvhA4GS0OdPVhTXteGA= -github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3/go.mod h1:3/Bu91CJONgkDA+Y0v/cnbROSJnu5tQ09vv7JGybUBA= +github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f h1:IYHVxTMV8oXlL4ZaoGDvREUZlszfyA0w3fUK9g2wSlM= +github.com/tikv/pd/client v0.0.0-20260729061341-533858f82c9f/go.mod h1:Xhd/CPuOv+tgP+5LKhWS23tklCXszF0EgKo4R6OZGj0= github.com/tjfoc/gmsm v1.4.1 h1:aMe1GlZb+0bLjn+cKTPEvvn9oUEBlJitaZiiBwsbgho= github.com/tjfoc/gmsm v1.4.1/go.mod h1:j4INPkHWMrhJb38G+J6W4Tw0AbuN8Thu3PbdVYhVcTE= github.com/tklauser/go-sysconf v0.3.9/go.mod h1:11DU/5sG7UexIrp/O6g35hrWzu0JxlwQ3LSFUzyeuhs= diff --git a/internal/apicodec/codec.go b/internal/apicodec/codec.go index b30bc550ea..fa4e0d2de4 100644 --- a/internal/apicodec/codec.go +++ b/internal/apicodec/codec.go @@ -6,6 +6,7 @@ import ( "github.com/pingcap/errors" "github.com/pingcap/kvproto/pkg/keyspacepb" "github.com/pingcap/kvproto/pkg/kvrpcpb" + "github.com/pingcap/kvproto/pkg/mpp" "github.com/tikv/client-go/v2/tikvrpc" "github.com/tikv/pd/client/constants" ) @@ -107,16 +108,28 @@ func DecodeKey(encoded []byte, version kvrpcpb.APIVersion) ([]byte, []byte, erro return nil, nil, err } return encoded[:keyspacePrefixLen], encoded[keyspacePrefixLen:], nil + case kvrpcpb.APIVersion_V3: + err := checkV3Key(encoded) + if err != nil { + return nil, nil, err + } + return encoded[:apiV3KeyspacePrefixLen], encoded[apiV3KeyspacePrefixLen:], nil } return nil, nil, errors.Errorf("unsupported api version %s", version.String()) } func setAPICtx(c Codec, r *tikvrpc.Request) { r.ApiVersion = c.GetAPIVersion() - r.KeyspaceId = uint32(c.GetKeyspaceID()) keyspaceMeta := c.GetKeyspaceMeta() if keyspaceMeta != nil { r.KeyspaceName = keyspaceMeta.Name + if identity := keyspaceMeta.GetKeyspaceIdentity(); identity != nil { + r.Keyspace = &kvrpcpb.Context_KeyspaceIdentity{KeyspaceIdentity: identity} + } else { + r.Keyspace = &kvrpcpb.Context_KeyspaceId{KeyspaceId: uint32(c.GetKeyspaceID())} + } + } else { + r.Keyspace = &kvrpcpb.Context_KeyspaceId{KeyspaceId: uint32(c.GetKeyspaceID())} } switch r.Type { @@ -124,15 +137,31 @@ func setAPICtx(c Codec, r *tikvrpc.Request) { mpp := *r.DispatchMPPTask() // Shallow copy the meta to avoid concurrent modification. meta := *mpp.Meta - meta.KeyspaceId = r.KeyspaceId meta.ApiVersion = r.ApiVersion + setMPPKeyspace(&meta, r) mpp.Meta = &meta r.Req = &mpp case tikvrpc.CmdCompact: compact := *r.Compact() - compact.KeyspaceId = r.KeyspaceId compact.ApiVersion = r.ApiVersion + setCompactKeyspace(&compact, r) r.Req = &compact } } + +func setMPPKeyspace(meta *mpp.TaskMeta, r *tikvrpc.Request) { + if identity := r.GetKeyspaceIdentity(); identity != nil { + meta.Keyspace = &mpp.TaskMeta_KeyspaceIdentity{KeyspaceIdentity: identity} + return + } + meta.Keyspace = &mpp.TaskMeta_KeyspaceId{KeyspaceId: r.GetKeyspaceId()} +} + +func setCompactKeyspace(compact *kvrpcpb.CompactRequest, r *tikvrpc.Request) { + if identity := r.GetKeyspaceIdentity(); identity != nil { + compact.Keyspace = &kvrpcpb.CompactRequest_KeyspaceIdentity{KeyspaceIdentity: identity} + return + } + compact.Keyspace = &kvrpcpb.CompactRequest_KeyspaceId{KeyspaceId: r.GetKeyspaceId()} +} diff --git a/internal/apicodec/codec_test.go b/internal/apicodec/codec_test.go index 3344b203ce..ae61577e42 100644 --- a/internal/apicodec/codec_test.go +++ b/internal/apicodec/codec_test.go @@ -47,6 +47,21 @@ func TestDecodeKey(t *testing.T) { assert.NotNil(t, err) assert.Empty(t, pfx) assert.Empty(t, key) + + pfx, key, err = DecodeKey([]byte{'x', 1, 2, 3, 4, 5, 6, 7, 8, 9}, kvrpcpb.APIVersion_V3) + assert.Nil(t, err) + assert.Equal(t, []byte{'x', 1, 2, 3}, pfx) + assert.Equal(t, []byte{4, 5, 6, 7, 8, 9}, key) + + pfx, key, err = DecodeKey([]byte{'x', 1, 2, 3}, kvrpcpb.APIVersion_V3) + assert.Nil(t, err) + assert.Equal(t, []byte{'x', 1, 2, 3}, pfx) + assert.Empty(t, key) + + pfx, key, err = DecodeKey([]byte{'x', 1, 2}, kvrpcpb.APIVersion_V3) + assert.NotNil(t, err) + assert.Empty(t, pfx) + assert.Empty(t, key) } func TestEncodeUnknownRequest(t *testing.T) { diff --git a/internal/apicodec/codec_v2.go b/internal/apicodec/codec_v2.go index 070e75fe91..2f71497140 100644 --- a/internal/apicodec/codec_v2.go +++ b/internal/apicodec/codec_v2.go @@ -55,6 +55,7 @@ func BuildKeyspaceName(name string) string { // codecV2 is used to encode/decode keys and request into APIv2 format. type codecV2 struct { reqPool sync.Pool + apiVersion kvrpcpb.APIVersion prefix []byte endKey []byte memCodec memCodec @@ -64,7 +65,7 @@ type codecV2 struct { // NewCodecV2 returns a codec that can be used to encode/decode // keys and requests to and from APIv2 format. func NewCodecV2(mode Mode, keyspaceMeta *keyspacepb.KeyspaceMeta) (Codec, error) { - keyspaceID := keyspaceMeta.Id + keyspaceID := keyspaceIDFromMeta(keyspaceMeta) if keyspaceID > maxKeyspaceID { return nil, errors.Errorf("keyspaceID %d is out of range, maximum is %d", keyspaceID, maxKeyspaceID) } @@ -74,6 +75,7 @@ func NewCodecV2(mode Mode, keyspaceMeta *keyspacepb.KeyspaceMeta) (Codec, error) } codec := &codecV2{ // Region keys in CodecV2 are always encoded in memory comparable form. + apiVersion: kvrpcpb.APIVersion_V2, memCodec: &memComparableCodec{}, keyspaceMeta: keyspaceMeta, } @@ -114,7 +116,7 @@ func (c *codecV2) GetKeyspace() []byte { } func (c *codecV2) GetKeyspaceID() KeyspaceID { - return KeyspaceID(c.keyspaceMeta.Id) + return KeyspaceID(keyspaceIDFromMeta(c.keyspaceMeta)) } func (c *codecV2) GetKeyspaceMeta() *keyspacepb.KeyspaceMeta { @@ -122,7 +124,14 @@ func (c *codecV2) GetKeyspaceMeta() *keyspacepb.KeyspaceMeta { } func (c *codecV2) GetAPIVersion() kvrpcpb.APIVersion { - return kvrpcpb.APIVersion_V2 + return c.apiVersion +} + +func keyspaceIDFromMeta(meta *keyspacepb.KeyspaceMeta) uint32 { + if identity := meta.GetKeyspaceIdentity(); identity != nil { + return identity.GetKeyspaceId() + } + return meta.GetId() } // EncodeRequest encodes with the given Codec. @@ -778,6 +787,10 @@ func (c *codecV2) encodeRange(start, end []byte, reverse bool) ([]byte, []byte) // DecodeRange maps encodedStart and end back to normal start and // end without APIv2 prefixes. func (c *codecV2) DecodeRange(encodedStart, encodedEnd []byte) (start []byte, end []byte, err error) { + if len(c.prefix) == 0 { + return encodedStart, encodedEnd, nil + } + if bytes.Compare(encodedStart, c.endKey) >= 0 || (len(encodedEnd) > 0 && bytes.Compare(encodedEnd, c.prefix) <= 0) { return nil, nil, errors.WithStack(errKeyOutOfBound) diff --git a/internal/apicodec/codec_v2_test.go b/internal/apicodec/codec_v2_test.go index f92ba59ee0..c40a0fda10 100644 --- a/internal/apicodec/codec_v2_test.go +++ b/internal/apicodec/codec_v2_test.go @@ -42,7 +42,7 @@ func TestCodecV2(t *testing.T) { func (suite *testCodecV2Suite) SetupSuite() { testKeyspaceMeta := keyspacepb.KeyspaceMeta{ - Id: testKeyspaceID, + Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: testKeyspaceID}, } codec, err := NewCodecV2(ModeRaw, &testKeyspaceMeta) suite.NoError(err) @@ -181,7 +181,7 @@ func (suite *testCodecV2Suite) TestEncodeRequest() { // TiFlash carries API v2 context separately through ApiVersion and KeyspaceId. re.Equal([]byte("compact-start"), encoded.Compact().StartKey) re.Equal(kvrpcpb.APIVersion_V2, encoded.Compact().ApiVersion) - re.Equal(testKeyspaceID, encoded.Compact().KeyspaceId) + re.Equal(testKeyspaceID, encoded.Compact().GetKeyspaceId()) }, }, } @@ -289,7 +289,7 @@ func (suite *testCodecV2Suite) TestNewCodecV2() { } for _, testCase := range testCases { - keyspaceMeta := &keyspacepb.KeyspaceMeta{Id: testCase.keyspaceID} + keyspaceMeta := &keyspacepb.KeyspaceMeta{Keyspace: &keyspacepb.KeyspaceMeta_Id{Id: testCase.keyspaceID}} if testCase.shouldErr { _, err := NewCodecV2(testCase.mode, keyspaceMeta) re.Error(err) @@ -884,7 +884,7 @@ func (suite *testCodecV2Suite) TestEncodeMPPRequest() { suite.Nil(err) task, ok := req.Req.(*mpp.DispatchTaskRequest) suite.True(ok) - suite.Equal(task.Meta.KeyspaceId, testKeyspaceID) + suite.Equal(testKeyspaceID, task.Meta.GetKeyspaceId()) suite.Equal(task.Meta.ApiVersion, kvrpcpb.APIVersion_V2) suite.Equal(task.Regions[0].Ranges[0].Start, suite.codec.EncodeKey([]byte("a"))) suite.Equal(task.Regions[0].Ranges[0].End, suite.codec.EncodeKey([]byte("b"))) diff --git a/internal/apicodec/codec_v3.go b/internal/apicodec/codec_v3.go new file mode 100644 index 0000000000..b4d3c87fdc --- /dev/null +++ b/internal/apicodec/codec_v3.go @@ -0,0 +1,225 @@ +package apicodec + +import ( + "bytes" + "encoding/binary" + + "github.com/pingcap/kvproto/pkg/apipb" + "github.com/pingcap/kvproto/pkg/keyspacepb" + "github.com/pingcap/kvproto/pkg/kvrpcpb" + "github.com/pkg/errors" + "github.com/tikv/client-go/v2/tikvrpc" +) + +const apiV3KeyspacePrefixLen = keyspacePrefixLen + +func checkV3Key(b []byte) error { + if len(b) < apiV3KeyspacePrefixLen || (b[0] != RawModePrefix && b[0] != TxnModePrefix) { + return errors.Errorf("invalid API V3 key %s", b) + } + return nil +} + +// codecV3 uses API V3 request context for TiKV RPCs. TiKV RPC keys stay +// logical and are scoped by the request context, while PD region lookups use +// the physical keyspace range so the region cache cannot cross into another +// API V3 keyspace. +type codecV3 struct { + *codecV2 + physicalPrefix []byte + physicalEndKey []byte +} + +// NewCodecV3 returns a codec for API V3 tenant-scoped keyspaces. +func NewCodecV3(mode Mode, identity *apipb.KeyspaceIdentity, keyspaceName string) (Codec, error) { + if identity == nil { + return nil, errors.New("missing API V3 keyspace identity") + } + namespaceID := identity.GetNamespaceId() + keyspaceID := identity.GetKeyspaceId() + if namespaceID == 0 { + return nil, errors.New("API V3 namespaceID must be non-zero") + } + if keyspaceID == 0 || keyspaceID > maxKeyspaceID { + return nil, errors.Errorf("API V3 keyspaceID %d is out of range, valid range is [1, %d]", keyspaceID, maxKeyspaceID) + } + + physicalPrefix := make([]byte, apiV3KeyspacePrefixLen) + switch mode { + case ModeRaw: + physicalPrefix[0] = RawModePrefix + case ModeTxn: + physicalPrefix[0] = TxnModePrefix + default: + return nil, errors.Errorf("unknown mode") + } + keyspaceIDBytes, err := getIDByte(keyspaceID) + if err != nil { + return nil, err + } + copy(physicalPrefix[1:], keyspaceIDBytes) + + physicalEndKey := make([]byte, apiV3KeyspacePrefixLen) + prefixVal := binary.BigEndian.Uint32(physicalPrefix) + binary.BigEndian.PutUint32(physicalEndKey, prefixVal+1) + + keyspaceMeta := &keyspacepb.KeyspaceMeta{ + Name: BuildKeyspaceName(keyspaceName), + Keyspace: &keyspacepb.KeyspaceMeta_KeyspaceIdentity{KeyspaceIdentity: identity}, + } + base := &codecV2{ + apiVersion: kvrpcpb.APIVersion_V3, + memCodec: &memComparableCodec{}, + keyspaceMeta: keyspaceMeta, + } + base.reqPool.New = func() any { return &tikvrpc.Request{} } + return &codecV3{codecV2: base, physicalPrefix: physicalPrefix, physicalEndKey: physicalEndKey}, nil +} + +func (c *codecV3) DecodeResponse(req *tikvrpc.Request, resp *tikvrpc.Response) (*tikvrpc.Response, error) { + return c.codecV2.DecodeResponse(req, resp) +} + +func (c *codecV3) GetKeyspace() []byte { + return c.physicalPrefix +} + +func (c *codecV3) EncodeRegionKey(key []byte) []byte { + return c.memCodec.encodeKey(c.encodePhysicalKey(key)) +} + +func (c *codecV3) DecodeRegionKey(encodedKey []byte) ([]byte, error) { + if len(encodedKey) == 0 { + return encodedKey, nil + } + key, err := c.memCodec.decodeKey(encodedKey) + if err != nil { + return nil, err + } + return c.decodeRegionKey(key) +} + +func (c *codecV3) EncodeRegionRange(start, end []byte) ([]byte, []byte) { + encodedEnd := c.physicalEndKey + if len(end) > 0 { + encodedEnd = c.encodePhysicalKey(end) + } + return c.memCodec.encodeKey(c.encodePhysicalKey(start)), c.memCodec.encodeKey(encodedEnd) +} + +func (c *codecV3) DecodeRegionRange(encodedStart, encodedEnd []byte) ([]byte, []byte, error) { + start, err := c.decodeRegionStart(encodedStart) + if err != nil { + return nil, nil, err + } + end, err := c.decodeRegionEnd(encodedEnd) + if err != nil { + return nil, nil, err + } + return start, end, nil +} + +func (c *codecV3) DecodeBucketKeys(keys [][]byte) ([][]byte, error) { + ks := make([][]byte, 0, len(keys)) + for i, key := range keys { + var ( + k []byte + err error + ) + if len(key) > 0 { + k, err = c.memCodec.decodeKey(key) + } + if err != nil { + return nil, err + } + if isV3PhysicalKey(k) { + if i == 0 && bytes.Compare(k, c.physicalPrefix) < 0 { + ks = append(ks, []byte{}) + } else if i == len(keys)-1 && (len(k) == 0 || bytes.Compare(k, c.physicalEndKey) >= 0) { + ks = append(ks, []byte{}) + } else if bytes.HasPrefix(k, c.physicalPrefix) { + raw := k[len(c.physicalPrefix):] + if len(raw) == 0 && len(ks) > 0 && len(ks[0]) == 0 { + continue + } + ks = append(ks, raw) + } + continue + } + if len(k) == 0 && i == len(keys)-1 { + ks = append(ks, []byte{}) + continue + } + if len(k) == 0 && i == 0 { + ks = append(ks, []byte{}) + continue + } + ks = append(ks, k) + } + return ks, nil +} + +func (c *codecV3) encodePhysicalKey(key []byte) []byte { + if bytes.HasPrefix(key, c.physicalPrefix) { + return key + } + encoded := make([]byte, 0, len(c.physicalPrefix)+len(key)) + encoded = append(encoded, c.physicalPrefix...) + encoded = append(encoded, key...) + return encoded +} + +func (c *codecV3) decodeRegionStart(encodedStart []byte) ([]byte, error) { + if len(encodedStart) == 0 { + return []byte{}, nil + } + start, err := c.memCodec.decodeKey(encodedStart) + if err != nil { + return nil, err + } + if isV3PhysicalKey(start) { + if bytes.Compare(start, c.physicalEndKey) >= 0 { + return nil, errors.WithStack(errKeyOutOfBound) + } + if bytes.Compare(start, c.physicalPrefix) < 0 { + return []byte{}, nil + } + } + return c.decodeRegionKey(start) +} + +func (c *codecV3) decodeRegionEnd(encodedEnd []byte) ([]byte, error) { + if len(encodedEnd) == 0 { + return []byte{}, nil + } + end, err := c.memCodec.decodeKey(encodedEnd) + if err != nil { + return nil, err + } + if isV3PhysicalKey(end) { + if bytes.Compare(end, c.physicalEndKey) >= 0 { + return []byte{}, nil + } + if bytes.Compare(end, c.physicalPrefix) <= 0 { + return nil, errors.WithStack(errKeyOutOfBound) + } + } + return c.decodeRegionKey(end) +} + +func (c *codecV3) decodeRegionKey(key []byte) ([]byte, error) { + if len(key) == 0 { + return []byte{}, nil + } + if bytes.HasPrefix(key, c.physicalPrefix) { + return key[len(c.physicalPrefix):], nil + } + if isV3PhysicalKey(key) { + return nil, errors.WithStack(errKeyOutOfBound) + } + return key, nil +} + +func isV3PhysicalKey(key []byte) bool { + return len(key) >= apiV3KeyspacePrefixLen && (key[0] == RawModePrefix || key[0] == TxnModePrefix) +} diff --git a/internal/apicodec/codec_v3_test.go b/internal/apicodec/codec_v3_test.go new file mode 100644 index 0000000000..7938a00a06 --- /dev/null +++ b/internal/apicodec/codec_v3_test.go @@ -0,0 +1,221 @@ +package apicodec + +import ( + "testing" + + "github.com/pingcap/kvproto/pkg/apipb" + "github.com/pingcap/kvproto/pkg/kvrpcpb" + "github.com/stretchr/testify/require" + "github.com/tikv/client-go/v2/tikvrpc" +) + +func TestNewCodecV3(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeRaw, identity, "ks") + re.NoError(err) + + v3Codec := codec.(*codecV3) + re.Equal(kvrpcpb.APIVersion_V3, v3Codec.GetAPIVersion()) + re.Equal([]byte{'r', 5, 6, 7}, v3Codec.physicalPrefix) + re.Equal([]byte{'r', 5, 6, 8}, v3Codec.physicalEndKey) + re.Equal(identity, v3Codec.GetKeyspaceMeta().GetKeyspaceIdentity()) + re.Equal(KeyspaceID(identity.KeyspaceId), v3Codec.GetKeyspaceID()) +} + +func TestNewCodecV3InvalidIdentity(t *testing.T) { + re := require.New(t) + + _, err := NewCodecV3(ModeTxn, nil, "") + re.Error(err) + + _, err = NewCodecV3(ModeTxn, &apipb.KeyspaceIdentity{NamespaceId: 0, KeyspaceId: 1}, "") + re.Error(err) + + _, err = NewCodecV3(ModeTxn, &apipb.KeyspaceIdentity{NamespaceId: 1, KeyspaceId: 0}, "") + re.Error(err) + + _, err = NewCodecV3(ModeTxn, &apipb.KeyspaceIdentity{NamespaceId: 1, KeyspaceId: maxKeyspaceID + 1}, "") + re.Error(err) +} + +func TestCodecV3EncodeRequestUsesLogicalKeys(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{NamespaceId: 7, KeyspaceId: 9} + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + + req := &tikvrpc.Request{ + Type: tikvrpc.CmdGet, + Req: &kvrpcpb.GetRequest{ + Key: []byte("key"), + }, + } + encoded, err := codec.EncodeRequest(req) + re.NoError(err) + defer codec.(*codecV3).reqPool.Put(encoded) + + re.Equal([]byte("key"), encoded.Get().Key) + re.Equal(kvrpcpb.APIVersion_V3, encoded.ApiVersion) + re.Equal("ks", encoded.KeyspaceName) + re.Equal(identity, encoded.GetKeyspaceIdentity()) +} + +func TestCodecV3EncodeRequestDoesNotEncodeScanBounds(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{NamespaceId: 7, KeyspaceId: 9} + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + + req := &tikvrpc.Request{ + Type: tikvrpc.CmdScan, + Req: &kvrpcpb.ScanRequest{ + StartKey: []byte("a"), + EndKey: []byte("b"), + }, + } + encoded, err := codec.EncodeRequest(req) + re.NoError(err) + defer codec.(*codecV3).reqPool.Put(encoded) + + re.Equal([]byte("a"), encoded.Scan().StartKey) + re.Equal([]byte("b"), encoded.Scan().EndKey) +} + +func TestCodecV3RegionKeysUsePhysicalPDKeys(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + + re.Equal([]byte("key"), codec.EncodeKey([]byte("key"))) + + regionKey := codec.EncodeRegionKey([]byte("key")) + physicalKey, err := codec.(*codecV3).memCodec.decodeKey(regionKey) + re.NoError(err) + re.Equal([]byte{'x', 5, 6, 7, 'k', 'e', 'y'}, physicalKey) + + decodedKey, err := codec.DecodeRegionKey(regionKey) + re.NoError(err) + re.Equal([]byte("key"), decodedKey) +} + +func TestCodecV3RegionRangeIsBoundedToPhysicalKeyspace(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + v3Codec := codec.(*codecV3) + + encodedStart, encodedEnd := codec.EncodeRegionRange([]byte("a"), nil) + + physicalStart, err := v3Codec.memCodec.decodeKey(encodedStart) + re.NoError(err) + re.Equal([]byte{'x', 5, 6, 7, 'a'}, physicalStart) + physicalEnd, err := v3Codec.memCodec.decodeKey(encodedEnd) + re.NoError(err) + re.Equal([]byte{'x', 5, 6, 8}, physicalEnd) + + decodedStart, decodedEnd, err := codec.DecodeRegionRange(encodedStart, encodedEnd) + re.NoError(err) + re.Equal([]byte("a"), decodedStart) + re.Empty(decodedEnd) +} + +func TestCodecV3DecodeRegionRangeAcceptsScopedLogicalKeys(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + v3Codec := codec.(*codecV3) + + encodedStart := v3Codec.memCodec.encodeKey([]byte("a")) + encodedEnd := v3Codec.memCodec.encodeKey([]byte("b")) + + decodedStart, decodedEnd, err := codec.DecodeRegionRange(encodedStart, encodedEnd) + re.NoError(err) + re.Equal([]byte("a"), decodedStart) + re.Equal([]byte("b"), decodedEnd) +} + +func TestCodecV3DecodeRegionRangeClampsPhysicalBounds(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + v3Codec := codec.(*codecV3) + + encodedStart := v3Codec.memCodec.encodeKey([]byte{'x', 5, 6, 6}) + encodedEnd := v3Codec.memCodec.encodeKey([]byte{'x', 5, 6, 8}) + + decodedStart, decodedEnd, err := codec.DecodeRegionRange(encodedStart, encodedEnd) + re.NoError(err) + re.Empty(decodedStart) + re.Empty(decodedEnd) +} + +func TestCodecV3DecodeRegionRangeRejectsOtherKeyspaceEnd(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + v3Codec := codec.(*codecV3) + + encodedStart := v3Codec.memCodec.encodeKey([]byte{'x', 5, 6, 7, 'a'}) + encodedEnd := v3Codec.memCodec.encodeKey([]byte{'x', 5, 6, 6, 'z'}) + + _, _, err = codec.DecodeRegionRange(encodedStart, encodedEnd) + re.Error(err) +} + +func TestCodecV3DecodeBucketKeysClampsPhysicalBounds(t *testing.T) { + re := require.New(t) + identity := &apipb.KeyspaceIdentity{ + NamespaceId: 0x01020304, + KeyspaceId: 0x050607, + } + codec, err := NewCodecV3(ModeTxn, identity, "ks") + re.NoError(err) + v3Codec := codec.(*codecV3) + + encodePhysical := func(key []byte) []byte { + return v3Codec.memCodec.encodeKey(key) + } + bucketKeys := [][]byte{ + encodePhysical([]byte{'x', 5, 6, 6, 'a'}), + codec.EncodeRegionKey([]byte{}), + codec.EncodeRegionKey([]byte("a")), + codec.EncodeRegionKey([]byte("b")), + codec.EncodeRegionKey([]byte("c")), + encodePhysical([]byte{'x', 5, 6, 8}), + encodePhysical([]byte{'x', 5, 6, 8, 'a'}), + } + + keys, err := codec.DecodeBucketKeys(bucketKeys) + re.NoError(err) + re.Equal([][]byte{ + {}, + []byte("a"), + []byte("b"), + []byte("c"), + {}, + }, keys) +} diff --git a/internal/locate/pd_codec.go b/internal/locate/pd_codec.go index 4e40a1f8ee..8689d3c822 100644 --- a/internal/locate/pd_codec.go +++ b/internal/locate/pd_codec.go @@ -37,6 +37,7 @@ package locate import ( "context" + "github.com/pingcap/kvproto/pkg/apipb" "github.com/pingcap/kvproto/pkg/keyspacepb" "github.com/pingcap/kvproto/pkg/pdpb" "github.com/pkg/errors" @@ -77,6 +78,16 @@ func NewCodecPDClientWithKeyspace(mode apicodec.Mode, client pd.Client, keyspace return &CodecPDClient{client.WithCallerComponent(componentName), codec}, nil } +// NewCodecPDClientWithKeyspaceIdentity creates a CodecPDClient in API v3 with keyspace identity. +func NewCodecPDClientWithKeyspaceIdentity(mode apicodec.Mode, client pd.Client, identity *apipb.KeyspaceIdentity, keyspace string) (*CodecPDClient, error) { + codec, err := apicodec.NewCodecV3(mode, identity, keyspace) + if err != nil { + return nil, err + } + + return &CodecPDClient{client.WithCallerComponent(componentName), codec}, nil +} + // GetKeyspaceID attempts to retrieve keyspace ID corresponding to the given keyspace name from PD. func GetKeyspaceID(client pd.Client, name string) (uint32, error) { meta, err := client.LoadKeyspace(context.Background(), apicodec.BuildKeyspaceName(name)) @@ -87,7 +98,7 @@ func GetKeyspaceID(client pd.Client, name string) (uint32, error) { if meta.State != keyspacepb.KeyspaceState_ENABLED { return 0, errors.Errorf("keyspace %s not enabled", name) } - return meta.Id, nil + return meta.GetId(), nil } // GetKeyspaceMeta attempts to retrieve keyspace meta corresponding to the given keyspace name from PD. diff --git a/internal/locate/region_cache.go b/internal/locate/region_cache.go index e2b69bde6c..cb39b8334d 100644 --- a/internal/locate/region_cache.go +++ b/internal/locate/region_cache.go @@ -816,6 +816,18 @@ func (c *RegionCache) Close() { c.bg.shutdown(true) } +func (c *RegionCache) allowRouterRegionLookup() bool { + return c.codec == nil || c.codec.GetAPIVersion() != kvrpcpb.APIVersion_V3 +} + +func (c *RegionCache) followerRegionOptions() []opt.GetRegionOption { + opts := []opt.GetRegionOption{opt.WithAllowFollowerHandle()} + if c.allowRouterRegionLookup() { + opts = append(opts, opt.WithAllowRouterServiceHandle()) + } + return opts +} + // IsBackgroundRunnerClosed returns whether RegionCache's background runner context has been canceled. // // It is intended for tests/debugging only. @@ -1783,7 +1795,7 @@ func (c *RegionCache) findRegionByKey(bo *retry.Backoffer, key []byte, isEndKey if r == nil || expired { // load region when it is not exists or expired. observeLoadRegion(tag, r, expired, 0) - lr, err := c.loadRegion(bo, key, isEndKey, opt.WithAllowFollowerHandle(), opt.WithAllowRouterServiceHandle()) + lr, err := c.loadRegion(bo, key, isEndKey, c.followerRegionOptions()...) if err != nil { // no region data, return error if failure. return nil, err @@ -2513,7 +2525,7 @@ func (c *RegionCache) scanRegions(bo *retry.Backoffer, startKey, endKey []byte, ctx = opentracing.ContextWithSpan(ctx, span1) } - pdOpts := []opt.GetRegionOption{opt.WithAllowFollowerHandle(), opt.WithAllowRouterServiceHandle()} + pdOpts := c.followerRegionOptions() var backoffErr error for { if backoffErr != nil { @@ -2598,7 +2610,7 @@ func (c *RegionCache) batchScanRegions(bo *retry.Backoffer, keyRanges []router.K opt.WithOutputMustContainAllKeyRange(), } if needFollowerHandle { - pdOpts = append(pdOpts, opt.WithAllowFollowerHandle(), opt.WithAllowRouterServiceHandle()) + pdOpts = append(pdOpts, c.followerRegionOptions()...) } if batchOpt.needBuckets { pdOpts = append(pdOpts, opt.WithBuckets()) diff --git a/internal/locate/region_cache_test.go b/internal/locate/region_cache_test.go index 4f1fcabe3c..9ba5732103 100644 --- a/internal/locate/region_cache_test.go +++ b/internal/locate/region_cache_test.go @@ -52,6 +52,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pingcap/failpoint" + "github.com/pingcap/kvproto/pkg/apipb" "github.com/pingcap/kvproto/pkg/errorpb" "github.com/pingcap/kvproto/pkg/kvrpcpb" "github.com/pingcap/kvproto/pkg/metapb" @@ -97,6 +98,31 @@ func (c *inspectedPDClient) BatchScanRegions(ctx context.Context, keyRanges []ro return c.Client.BatchScanRegions(ctx, keyRanges, limit, opts...) } +func TestFollowerRegionOptionsDisableRouterForAPIV3(t *testing.T) { + t.Parallel() + + v1 := &RegionCache{codec: apicodec.NewCodecV1(apicodec.ModeTxn)} + var v1Op opt.GetRegionOp + for _, o := range v1.followerRegionOptions() { + o(&v1Op) + } + require.True(t, v1Op.AllowFollowerHandle) + require.True(t, v1Op.AllowRouterServiceHandle) + + v3Codec, err := apicodec.NewCodecV3(apicodec.ModeTxn, &apipb.KeyspaceIdentity{ + NamespaceId: 7, + KeyspaceId: 11, + }, "ks") + require.NoError(t, err) + v3 := &RegionCache{codec: v3Codec} + var v3Op opt.GetRegionOp + for _, o := range v3.followerRegionOptions() { + o(&v3Op) + } + require.True(t, v3Op.AllowFollowerHandle) + require.False(t, v3Op.AllowRouterServiceHandle) +} + func TestBackgroundRunner(t *testing.T) { t.Run("ShutdownWait", func(t *testing.T) { dur := 100 * time.Millisecond diff --git a/internal/mockstore/mocktikv/cluster.go b/internal/mockstore/mocktikv/cluster.go index 72a0875abb..12eb1d2a4e 100644 --- a/internal/mockstore/mocktikv/cluster.go +++ b/internal/mockstore/mocktikv/cluster.go @@ -196,10 +196,16 @@ func (c *Cluster) UnCancelStore(storeID uint64) { // GetStoreByAddr returns a Store's meta by an addr. func (c *Cluster) GetStoreByAddr(addr string) *metapb.Store { + if c == nil { + return nil + } c.RLock() defer c.RUnlock() for _, s := range c.stores { + if s == nil || s.meta == nil { + continue + } if s.meta.GetAddress() == addr { return proto.Clone(s.meta).(*metapb.Store) } @@ -209,10 +215,16 @@ func (c *Cluster) GetStoreByAddr(addr string) *metapb.Store { // GetAndCheckStoreByAddr checks and returns a Store's meta by an addr func (c *Cluster) GetAndCheckStoreByAddr(addr string) (ss []*metapb.Store, err error) { + if c == nil { + return nil, context.Canceled + } c.RLock() defer c.RUnlock() for _, s := range c.stores { + if s == nil || s.meta == nil { + continue + } if s.cancel { err = context.Canceled return diff --git a/tikv/compatible_txn_safe_point_loader.go b/tikv/compatible_txn_safe_point_loader.go index 013a7e38c2..a5b758dbec 100644 --- a/tikv/compatible_txn_safe_point_loader.go +++ b/tikv/compatible_txn_safe_point_loader.go @@ -104,7 +104,7 @@ func (l *compatibleTxnSafePointLoader) loadTxnSafePoint(ctx context.Context) (ui key := unifiedTxnSafePointPath keyspaceMeta := l.codec.GetKeyspaceMeta() if pd.IsKeyspaceUsingKeyspaceLevelGC(keyspaceMeta) { - key = fmt.Sprintf(keyspaceLevelTxnSafePointPath, keyspaceMeta.Id) + key = fmt.Sprintf(keyspaceLevelTxnSafePointPath, keyspaceMeta.GetId()) } // Follow the same implementation as the EtcdSafePointKV by setting the timeout 5 seconds. diff --git a/tikv/kv.go b/tikv/kv.go index 984093443a..1c3cf0676a 100644 --- a/tikv/kv.go +++ b/tikv/kv.go @@ -47,6 +47,7 @@ import ( "time" "github.com/opentracing/opentracing-go" + "github.com/pingcap/kvproto/pkg/apipb" "github.com/pingcap/kvproto/pkg/kvrpcpb" "github.com/pkg/errors" "github.com/tikv/client-go/v2/config" @@ -387,15 +388,17 @@ func NewKVStore( // NewPDClient returns an unwrapped pd client. func NewPDClient(pdAddrs []string) (pd.Client, error) { + return newPDClient(pdAddrs, nil) +} + +// NewPDClientWithKeyspaceIdentity returns an unwrapped API V3 pd client. +func NewPDClientWithKeyspaceIdentity(pdAddrs []string, identity *apipb.KeyspaceIdentity) (pd.Client, error) { + return newPDClient(pdAddrs, identity) +} + +func buildPDClientOptions(identity *apipb.KeyspaceIdentity) []opt.ClientOption { cfg := config.GetGlobalConfig() - // init pd-client - pdCli, err := pd.NewClient( - caller.Component("client-go"), - pdAddrs, pd.SecurityOption{ - CAPath: cfg.Security.ClusterSSLCA, - CertPath: cfg.Security.ClusterSSLCert, - KeyPath: cfg.Security.ClusterSSLKey, - }, + opts := []opt.ClientOption{ opt.WithGRPCDialOptions( grpc.WithKeepaliveParams( keepalive.ClientParameters{ @@ -404,9 +407,45 @@ func NewPDClient(pdAddrs []string) (pd.Client, error) { }, ), ), - opt.WithCustomTimeoutOption(time.Duration(cfg.PDClient.PDServerTimeout)*time.Second), + opt.WithCustomTimeoutOption(time.Duration(cfg.PDClient.PDServerTimeout) * time.Second), opt.WithForwardingOption(config.GetGlobalConfig().EnableForwarding), + } + if identity != nil { + opts = append(opts, opt.WithEnableRouterClient(false)) + } + return opts +} + +func newPDClient(pdAddrs []string, identity *apipb.KeyspaceIdentity) (pd.Client, error) { + cfg := config.GetGlobalConfig() + security := pd.SecurityOption{ + CAPath: cfg.Security.ClusterSSLCA, + CertPath: cfg.Security.ClusterSSLCert, + KeyPath: cfg.Security.ClusterSSLKey, + } + opts := buildPDClientOptions(identity) + + var ( + pdCli pd.Client + err error ) + if identity != nil { + pdCli, err = pd.NewClientWithKeyspaceIdentity( + context.Background(), + caller.Component("client-go"), + identity, + pdAddrs, + security, + opts..., + ) + } else { + pdCli, err = pd.NewClient( + caller.Component("client-go"), + pdAddrs, + security, + opts..., + ) + } if err != nil { return nil, errors.WithStack(err) } diff --git a/tikv/kv_test.go b/tikv/kv_test.go index 8a7042819d..b2612e9799 100644 --- a/tikv/kv_test.go +++ b/tikv/kv_test.go @@ -23,6 +23,7 @@ import ( "time" "github.com/pingcap/failpoint" + "github.com/pingcap/kvproto/pkg/apipb" "github.com/pingcap/kvproto/pkg/kvrpcpb" "github.com/pingcap/kvproto/pkg/metapb" "github.com/stretchr/testify/require" @@ -33,6 +34,7 @@ import ( "github.com/tikv/client-go/v2/tikvrpc" "github.com/tikv/client-go/v2/util" pdhttp "github.com/tikv/pd/client/http" + pdopt "github.com/tikv/pd/client/opt" ) func TestKV(t *testing.T) { @@ -40,6 +42,25 @@ func TestKV(t *testing.T) { suite.Run(t, new(testKVSuite)) } +func TestBuildPDClientOptionsDisableRouterForAPIV3(t *testing.T) { + t.Parallel() + + defaultOptions := pdopt.NewOption() + for _, o := range buildPDClientOptions(nil) { + o(defaultOptions) + } + require.True(t, defaultOptions.GetEnableRouterClient()) + + apiV3Options := pdopt.NewOption() + for _, o := range buildPDClientOptions(&apipb.KeyspaceIdentity{ + NamespaceId: 1, + KeyspaceId: 2, + }) { + o(apiV3Options) + } + require.False(t, apiV3Options.GetEnableRouterClient()) +} + type testKVSuite struct { suite.Suite store *KVStore diff --git a/tikv/region.go b/tikv/region.go index 309cd98f4d..21be840606 100644 --- a/tikv/region.go +++ b/tikv/region.go @@ -118,16 +118,22 @@ var NewCodecPDClient = locate.NewCodecPDClient // NewCodecPDClientWithKeyspace creates a CodecPDClient in API v2 with keyspace name. var NewCodecPDClientWithKeyspace = locate.NewCodecPDClientWithKeyspace +// NewCodecPDClientWithKeyspaceIdentity creates a CodecPDClient in API v3 with keyspace identity. +var NewCodecPDClientWithKeyspaceIdentity = locate.NewCodecPDClientWithKeyspaceIdentity + // NewCodecV1 is a constructor for v1 Codec. var NewCodecV1 = apicodec.NewCodecV1 // NewCodecV2 is a constructor for v2 Codec. var NewCodecV2 = apicodec.NewCodecV2 +// NewCodecV3 is a constructor for v3 Codec. +var NewCodecV3 = apicodec.NewCodecV3 + // Codec is responsible for encode/decode requests. type Codec = apicodec.Codec -// DecodeKey is used to split a given key to it's APIv2 prefix and actual key. +// DecodeKey is used to split a given key to its API keyspace prefix and actual key. var DecodeKey = apicodec.DecodeKey // DefaultKeyspaceID is the keyspaceID of the default keyspace. diff --git a/tikv/test_util.go b/tikv/test_util.go index ea4f4b5e04..786e3cfda4 100644 --- a/tikv/test_util.go +++ b/tikv/test_util.go @@ -140,7 +140,7 @@ func NewTestKeyspaceTiKVStore(client Client, pdClient pd.Client, clientHijack fu // Make sure the uuid is unique. uid := uuid.New().String() - keyspaceIdStr := strconv.FormatUint(uint64(keyspaceMeta.Id), 10) + keyspaceIdStr := strconv.FormatUint(uint64(keyspaceMeta.GetId()), 10) spkv := NewMockSafePointKV(WithPrefix(keyspaceIdStr)) tikvStore, err := NewKVStore(uid, pdCli, spkv, client, opt...) diff --git a/tikvrpc/tikvrpc_test.go b/tikvrpc/tikvrpc_test.go index d198fd9935..e1a4c1f0e3 100644 --- a/tikvrpc/tikvrpc_test.go +++ b/tikvrpc/tikvrpc_test.go @@ -125,13 +125,13 @@ func TestAttachContextSetsRequestContext(t *testing.T) { rpcCtx := kvrpcpb.Context{ RegionId: 123, ApiVersion: kvrpcpb.APIVersion_V2, - KeyspaceId: 456, + Keyspace: &kvrpcpb.Context_KeyspaceId{KeyspaceId: 456}, KeyspaceName: "test-keyspace", } nextRPCtx := kvrpcpb.Context{ RegionId: 789, ApiVersion: kvrpcpb.APIVersion_V2, - KeyspaceId: 101112, + Keyspace: &kvrpcpb.Context_KeyspaceId{KeyspaceId: 101112}, KeyspaceName: "next-test-keyspace", } diff --git a/txnkv/client.go b/txnkv/client.go index bfcaad47a7..22dcd73bbc 100644 --- a/txnkv/client.go +++ b/txnkv/client.go @@ -18,6 +18,7 @@ import ( "context" "fmt" + "github.com/pingcap/kvproto/pkg/apipb" "github.com/pingcap/kvproto/pkg/kvrpcpb" "github.com/pkg/errors" "github.com/tikv/client-go/v2/config" @@ -26,6 +27,7 @@ import ( "github.com/tikv/client-go/v2/tikv" "github.com/tikv/client-go/v2/txnkv/transaction" "github.com/tikv/client-go/v2/util" + pd "github.com/tikv/pd/client" ) // Client is a txn client. @@ -34,9 +36,10 @@ type Client struct { } type option struct { - apiVersion kvrpcpb.APIVersion - keyspaceName string - spKVPrefix string + apiVersion kvrpcpb.APIVersion + keyspaceName string + keyspaceIdentity *apipb.KeyspaceIdentity + spKVPrefix string } // ClientOpt is factory to set the client options. @@ -49,6 +52,16 @@ func WithKeyspace(keyspaceName string) ClientOpt { } } +// WithKeyspaceIdentity is used to set client's API V3 keyspace identity. +func WithKeyspaceIdentity(namespaceID, keyspaceID uint32) ClientOpt { + return func(opt *option) { + opt.keyspaceIdentity = &apipb.KeyspaceIdentity{ + NamespaceId: namespaceID, + KeyspaceId: keyspaceID, + } + } +} + // WithAPIVersion is used to set client's apiVersion. func WithAPIVersion(apiVersion kvrpcpb.APIVersion) ClientOpt { return func(opt *option) { @@ -71,7 +84,13 @@ func NewClient(pdAddrs []string, opts ...ClientOpt) (*Client, error) { o(opt) } // Use an unwrapped PDClient to obtain keyspace meta. - pdClient, err := tikv.NewPDClient(pdAddrs) + newPDClient := tikv.NewPDClient + if opt.apiVersion == kvrpcpb.APIVersion_V3 { + newPDClient = func(pdAddrs []string) (pd.Client, error) { + return tikv.NewPDClientWithKeyspaceIdentity(pdAddrs, opt.keyspaceIdentity) + } + } + pdClient, err := newPDClient(pdAddrs) if err != nil { return nil, errors.WithStack(err) } @@ -88,6 +107,11 @@ func NewClient(pdAddrs []string, opts ...ClientOpt) (*Client, error) { if err != nil { return nil, err } + case kvrpcpb.APIVersion_V3: + codecCli, err = tikv.NewCodecPDClientWithKeyspaceIdentity(tikv.ModeTxn, pdClient, opt.keyspaceIdentity, opt.keyspaceName) + if err != nil { + return nil, err + } default: return nil, errors.Errorf("unknown api version: %d", opt.apiVersion) }