From ed6bd8ee251a2be274ec95d02fc72661c7ea6147 Mon Sep 17 00:00:00 2001 From: premal Date: Tue, 16 Jun 2026 21:04:01 +0000 Subject: [PATCH 1/2] feat: add encoding capabilities package for TiFlash dictionary encoding Adds PD-side tracking of TiFlash encoding capabilities: - StoreEncodingCapabilities: tracks per-store encoding support (dictionary encoding, encoded filter/group-by/bloom filter/star join) - ClusterEncodingStatus: cluster-wide view with AllStoresReady flag for safe dictionary encoding enablement - Registry: thread-safe store for encoding capability reporting - Handler: HTTP API endpoints for encoding status and registration (GET /pd/api/v1/encoding/status, POST /pd/api/v1/encoding/capabilities) - ComputeClusterStatus: computes min encoding version across stores - Comprehensive unit tests for all functionality This enables TiDB's planner to query PD for cluster encoding readiness before enabling dictionary-encoded operations on TiFlash nodes. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/encoded/capabilities.go | 139 +++++++++++++++++++++++++++++++ pkg/encoded/capabilities_test.go | 104 +++++++++++++++++++++++ pkg/encoded/handler.go | 114 +++++++++++++++++++++++++ 3 files changed, 357 insertions(+) create mode 100644 pkg/encoded/capabilities.go create mode 100644 pkg/encoded/capabilities_test.go create mode 100644 pkg/encoded/handler.go diff --git a/pkg/encoded/capabilities.go b/pkg/encoded/capabilities.go new file mode 100644 index 0000000000..50fcca795f --- /dev/null +++ b/pkg/encoded/capabilities.go @@ -0,0 +1,139 @@ +// Copyright 2024 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 encoded provides encoding capability metadata for TiFlash stores. +// This is part of the dictionary-encoded columnstore operations feature (Phase 1). +// +// When TiFlash nodes support dictionary encoding, they report their capabilities +// via the store heartbeat. PD tracks which stores support encoded operations, +// enabling TiDB's planner to make informed decisions about which TiFlash nodes +// can execute encoded operations. +package encoded + +// Capability represents a single encoding capability supported by a store. +type Capability string + +const ( + // CapDictionaryEncoding indicates the store supports dictionary encoding + // for low-cardinality columns in the columnar storage format. + CapDictionaryEncoding Capability = "dictionary_encoding" + + // CapEncodedFilter indicates the store can perform filter operations + // directly on dictionary-encoded data without decoding. + CapEncodedFilter Capability = "encoded_filter" + + // CapEncodedGroupBy indicates the store can perform group-by operations + // using dictionary IDs as direct array indices. + CapEncodedGroupBy Capability = "encoded_group_by" + + // CapEncodedBloomFilter indicates the store supports bloom filter pushdown + // on dictionary-encoded columns. + CapEncodedBloomFilter Capability = "encoded_bloom_filter" + + // CapEncodedStarJoin indicates the store supports encoded star join + // operations (hash join on encoded fact table columns). + CapEncodedStarJoin Capability = "encoded_star_join" +) + +// StoreEncodingCapabilities tracks what encoding operations a TiFlash store supports. +type StoreEncodingCapabilities struct { + // StoreID is the PD store ID for the TiFlash node. + StoreID uint64 `json:"store_id"` + + // Capabilities lists the encoding operations this store supports. + Capabilities []Capability `json:"capabilities"` + + // MaxDictCardinality is the maximum dictionary size this store handles. + MaxDictCardinality uint32 `json:"max_dict_cardinality"` + + // EncodingVersion indicates the protocol version for encoded operations. + // Version 1: Dictionary encoding (Phase 1) + // Version 2: Encoded filter + group-by (Phase 2-3) + // Version 3: Encoded star joins (Phase 4) + EncodingVersion uint32 `json:"encoding_version"` +} + +// HasCapability checks if a store supports a given encoding capability. +func (s *StoreEncodingCapabilities) HasCapability(cap Capability) bool { + for _, c := range s.Capabilities { + if c == cap { + return true + } + } + return false +} + +// SupportsEncodedOperations returns true if the store has any encoded operation +// capabilities beyond basic dictionary encoding. +func (s *StoreEncodingCapabilities) SupportsEncodedOperations() bool { + return s.HasCapability(CapEncodedFilter) || + s.HasCapability(CapEncodedGroupBy) || + s.HasCapability(CapEncodedBloomFilter) || + s.HasCapability(CapEncodedStarJoin) +} + +// NewStoreEncodingCapabilities creates capabilities for a store that supports +// Phase 1 dictionary encoding. +func NewStoreEncodingCapabilities(storeID uint64) *StoreEncodingCapabilities { + return &StoreEncodingCapabilities{ + StoreID: storeID, + Capabilities: []Capability{CapDictionaryEncoding}, + MaxDictCardinality: 4096, + EncodingVersion: 1, + } +} + +// ClusterEncodingStatus provides a cluster-wide view of encoding capabilities. +type ClusterEncodingStatus struct { + // TotalTiFlashStores is the number of TiFlash stores in the cluster. + TotalTiFlashStores int `json:"total_tiflash_stores"` + + // EncodingEnabledStores is the number of TiFlash stores that support + // dictionary encoding. + EncodingEnabledStores int `json:"encoding_enabled_stores"` + + // AllStoresReady is true when all TiFlash stores support dictionary encoding. + // This is the precondition for enabling dictionary encoding writes. + AllStoresReady bool `json:"all_stores_ready"` + + // MinEncodingVersion is the minimum encoding version across all stores. + // TiDB should only enable features up to this version. + MinEncodingVersion uint32 `json:"min_encoding_version"` + + // Stores contains per-store capability details. + Stores []*StoreEncodingCapabilities `json:"stores,omitempty"` +} + +// ComputeClusterStatus computes the cluster-wide encoding status from +// individual store capabilities. +func ComputeClusterStatus(stores []*StoreEncodingCapabilities, totalTiFlash int) *ClusterEncodingStatus { + status := &ClusterEncodingStatus{ + TotalTiFlashStores: totalTiFlash, + EncodingEnabledStores: len(stores), + AllStoresReady: len(stores) == totalTiFlash && totalTiFlash > 0, + MinEncodingVersion: 0, + Stores: stores, + } + + if len(stores) > 0 { + status.MinEncodingVersion = stores[0].EncodingVersion + for _, s := range stores[1:] { + if s.EncodingVersion < status.MinEncodingVersion { + status.MinEncodingVersion = s.EncodingVersion + } + } + } + + return status +} diff --git a/pkg/encoded/capabilities_test.go b/pkg/encoded/capabilities_test.go new file mode 100644 index 0000000000..5cfe60f442 --- /dev/null +++ b/pkg/encoded/capabilities_test.go @@ -0,0 +1,104 @@ +// Copyright 2024 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 encoded + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestStoreEncodingCapabilities(t *testing.T) { + re := require.New(t) + + // Test NewStoreEncodingCapabilities defaults + caps := NewStoreEncodingCapabilities(1) + re.Equal(uint64(1), caps.StoreID) + re.Equal(uint32(4096), caps.MaxDictCardinality) + re.Equal(uint32(1), caps.EncodingVersion) + re.True(caps.HasCapability(CapDictionaryEncoding)) + re.False(caps.HasCapability(CapEncodedFilter)) + re.False(caps.HasCapability(CapEncodedGroupBy)) + re.False(caps.SupportsEncodedOperations()) +} + +func TestStoreEncodingCapabilitiesWithOps(t *testing.T) { + re := require.New(t) + + caps := &StoreEncodingCapabilities{ + StoreID: 2, + Capabilities: []Capability{CapDictionaryEncoding, CapEncodedFilter, CapEncodedGroupBy}, + MaxDictCardinality: 8192, + EncodingVersion: 2, + } + re.True(caps.HasCapability(CapDictionaryEncoding)) + re.True(caps.HasCapability(CapEncodedFilter)) + re.True(caps.HasCapability(CapEncodedGroupBy)) + re.False(caps.HasCapability(CapEncodedStarJoin)) + re.True(caps.SupportsEncodedOperations()) +} + +func TestClusterEncodingStatus(t *testing.T) { + re := require.New(t) + + // All stores ready + stores := []*StoreEncodingCapabilities{ + NewStoreEncodingCapabilities(1), + NewStoreEncodingCapabilities(2), + NewStoreEncodingCapabilities(3), + } + status := ComputeClusterStatus(stores, 3) + re.True(status.AllStoresReady) + re.Equal(3, status.EncodingEnabledStores) + re.Equal(3, status.TotalTiFlashStores) + re.Equal(uint32(1), status.MinEncodingVersion) +} + +func TestClusterEncodingStatusPartial(t *testing.T) { + re := require.New(t) + + // Only 2 of 3 stores report encoding capabilities + stores := []*StoreEncodingCapabilities{ + NewStoreEncodingCapabilities(1), + NewStoreEncodingCapabilities(2), + } + status := ComputeClusterStatus(stores, 3) + re.False(status.AllStoresReady) + re.Equal(2, status.EncodingEnabledStores) + re.Equal(3, status.TotalTiFlashStores) +} + +func TestClusterEncodingStatusMixedVersions(t *testing.T) { + re := require.New(t) + + stores := []*StoreEncodingCapabilities{ + {StoreID: 1, Capabilities: []Capability{CapDictionaryEncoding}, EncodingVersion: 2}, + {StoreID: 2, Capabilities: []Capability{CapDictionaryEncoding}, EncodingVersion: 1}, + {StoreID: 3, Capabilities: []Capability{CapDictionaryEncoding}, EncodingVersion: 3}, + } + status := ComputeClusterStatus(stores, 3) + re.True(status.AllStoresReady) + // MinEncodingVersion should be the lowest + re.Equal(uint32(1), status.MinEncodingVersion) +} + +func TestClusterEncodingStatusEmpty(t *testing.T) { + re := require.New(t) + + status := ComputeClusterStatus(nil, 0) + re.False(status.AllStoresReady) + re.Equal(0, status.EncodingEnabledStores) + re.Equal(uint32(0), status.MinEncodingVersion) +} diff --git a/pkg/encoded/handler.go b/pkg/encoded/handler.go new file mode 100644 index 0000000000..f697f52a1c --- /dev/null +++ b/pkg/encoded/handler.go @@ -0,0 +1,114 @@ +// Copyright 2024 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 encoded + +import ( + "encoding/json" + "net/http" + "sync" +) + +// Registry tracks encoding capabilities reported by TiFlash stores. +// It is safe for concurrent access. +type Registry struct { + mu sync.RWMutex + store map[uint64]*StoreEncodingCapabilities +} + +// NewRegistry creates an empty encoding capabilities registry. +func NewRegistry() *Registry { + return &Registry{ + store: make(map[uint64]*StoreEncodingCapabilities), + } +} + +// Register adds or updates encoding capabilities for a store. +func (r *Registry) Register(caps *StoreEncodingCapabilities) { + r.mu.Lock() + defer r.mu.Unlock() + r.store[caps.StoreID] = caps +} + +// Unregister removes encoding capabilities for a store. +func (r *Registry) Unregister(storeID uint64) { + r.mu.Lock() + defer r.mu.Unlock() + delete(r.store, storeID) +} + +// Get retrieves encoding capabilities for a specific store. +func (r *Registry) Get(storeID uint64) (*StoreEncodingCapabilities, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + caps, ok := r.store[storeID] + return caps, ok +} + +// GetAll returns all registered store capabilities. +func (r *Registry) GetAll() []*StoreEncodingCapabilities { + r.mu.RLock() + defer r.mu.RUnlock() + result := make([]*StoreEncodingCapabilities, 0, len(r.store)) + for _, caps := range r.store { + result = append(result, caps) + } + return result +} + +// GetClusterStatus computes the cluster-wide encoding status. +func (r *Registry) GetClusterStatus(totalTiFlashStores int) *ClusterEncodingStatus { + stores := r.GetAll() + return ComputeClusterStatus(stores, totalTiFlashStores) +} + +// Handler provides HTTP API endpoints for encoding capabilities. +type Handler struct { + registry *Registry + getTiFlashCount func() int +} + +// NewHandler creates an API handler for encoding capabilities. +func NewHandler(registry *Registry, getTiFlashCount func() int) *Handler { + return &Handler{ + registry: registry, + getTiFlashCount: getTiFlashCount, + } +} + +// ServeEncodingStatus handles GET /pd/api/v1/encoding/status +func (h *Handler) ServeEncodingStatus(w http.ResponseWriter, r *http.Request) { + status := h.registry.GetClusterStatus(h.getTiFlashCount()) + + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + json.NewEncoder(w).Encode(status) +} + +// ServeRegisterCapabilities handles POST /pd/api/v1/encoding/capabilities +func (h *Handler) ServeRegisterCapabilities(w http.ResponseWriter, r *http.Request) { + var caps StoreEncodingCapabilities + if err := json.NewDecoder(r.Body).Decode(&caps); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + if caps.StoreID == 0 { + http.Error(w, "store_id is required", http.StatusBadRequest) + return + } + + h.registry.Register(&caps) + w.WriteHeader(http.StatusOK) +} From 663e3e7d131ee3decb0663597aea8834de8854ad Mon Sep 17 00:00:00 2001 From: premal Date: Tue, 16 Jun 2026 21:49:32 +0000 Subject: [PATCH 2/2] encoded: add Phase 2-4 capability tracking and feasibility checks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds: - NewFullEncodingCapabilities(): creates v3 store with all 5 capabilities - PhaseCapabilities map: phase→capabilities mapping - SupportsPhase()/ClusterSupportsPhase(): phase version checks - OperationFeasibility struct: per-operation cluster readiness - CheckOperationFeasibility()/AllOperationsFeasibility(): checks if all TiFlash stores support specific encoded operations 13 tests covering full capability lifecycle. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/encoded/capabilities.go | 85 +++++++++++++++++++++++++ pkg/encoded/capabilities_test.go | 105 +++++++++++++++++++++++++++++++ 2 files changed, 190 insertions(+) diff --git a/pkg/encoded/capabilities.go b/pkg/encoded/capabilities.go index 50fcca795f..01eb09b60f 100644 --- a/pkg/encoded/capabilities.go +++ b/pkg/encoded/capabilities.go @@ -137,3 +137,88 @@ func ComputeClusterStatus(stores []*StoreEncodingCapabilities, totalTiFlash int) return status } + +// NewFullEncodingCapabilities creates capabilities for a store that supports +// all 4 phases of encoded operations. +func NewFullEncodingCapabilities(storeID uint64) *StoreEncodingCapabilities { + return &StoreEncodingCapabilities{ + StoreID: storeID, + Capabilities: []Capability{ + CapDictionaryEncoding, + CapEncodedFilter, + CapEncodedGroupBy, + CapEncodedBloomFilter, + CapEncodedStarJoin, + }, + MaxDictCardinality: 4096, + EncodingVersion: 3, + } +} + +// PhaseCapabilities maps encoding versions to their supported capabilities. +var PhaseCapabilities = map[uint32][]Capability{ + 1: {CapDictionaryEncoding}, + 2: {CapDictionaryEncoding, CapEncodedFilter, CapEncodedGroupBy}, + 3: {CapDictionaryEncoding, CapEncodedFilter, CapEncodedGroupBy, CapEncodedBloomFilter, CapEncodedStarJoin}, +} + +// SupportsPhase checks if a store supports a specific encoding phase. +func (s *StoreEncodingCapabilities) SupportsPhase(phase uint32) bool { + return s.EncodingVersion >= phase +} + +// ClusterSupportsPhase checks if all stores in the cluster support a given phase. +func (status *ClusterEncodingStatus) ClusterSupportsPhase(phase uint32) bool { + return status.AllStoresReady && status.MinEncodingVersion >= phase +} + +// OperationFeasibility describes whether a specific encoded operation can run +// across the cluster. +type OperationFeasibility struct { + // Operation is the capability being checked + Operation Capability `json:"operation"` + + // Feasible indicates all relevant stores support this operation + Feasible bool `json:"feasible"` + + // StoresReady is the count of stores that support this operation + StoresReady int `json:"stores_ready"` + + // StoresTotal is the total number of TiFlash stores + StoresTotal int `json:"stores_total"` +} + +// CheckOperationFeasibility checks if a specific encoded operation can run +// on all TiFlash stores in the cluster. +func CheckOperationFeasibility(stores []*StoreEncodingCapabilities, totalTiFlash int, op Capability) *OperationFeasibility { + result := &OperationFeasibility{ + Operation: op, + StoresTotal: totalTiFlash, + } + + for _, s := range stores { + if s.HasCapability(op) { + result.StoresReady++ + } + } + + result.Feasible = result.StoresReady == totalTiFlash && totalTiFlash > 0 + return result +} + +// AllOperationsFeasibility checks feasibility for all encoded operations. +func AllOperationsFeasibility(stores []*StoreEncodingCapabilities, totalTiFlash int) []*OperationFeasibility { + ops := []Capability{ + CapDictionaryEncoding, + CapEncodedFilter, + CapEncodedGroupBy, + CapEncodedBloomFilter, + CapEncodedStarJoin, + } + + results := make([]*OperationFeasibility, 0, len(ops)) + for _, op := range ops { + results = append(results, CheckOperationFeasibility(stores, totalTiFlash, op)) + } + return results +} diff --git a/pkg/encoded/capabilities_test.go b/pkg/encoded/capabilities_test.go index 5cfe60f442..1ac229873e 100644 --- a/pkg/encoded/capabilities_test.go +++ b/pkg/encoded/capabilities_test.go @@ -102,3 +102,108 @@ func TestClusterEncodingStatusEmpty(t *testing.T) { re.Equal(0, status.EncodingEnabledStores) re.Equal(uint32(0), status.MinEncodingVersion) } + +func TestNewFullEncodingCapabilities(t *testing.T) { + re := require.New(t) + + caps := NewFullEncodingCapabilities(5) + re.Equal(uint64(5), caps.StoreID) + re.Equal(uint32(3), caps.EncodingVersion) + re.True(caps.HasCapability(CapDictionaryEncoding)) + re.True(caps.HasCapability(CapEncodedFilter)) + re.True(caps.HasCapability(CapEncodedGroupBy)) + re.True(caps.HasCapability(CapEncodedBloomFilter)) + re.True(caps.HasCapability(CapEncodedStarJoin)) + re.True(caps.SupportsEncodedOperations()) +} + +func TestSupportsPhase(t *testing.T) { + re := require.New(t) + + caps := NewStoreEncodingCapabilities(1) // Version 1 + re.True(caps.SupportsPhase(1)) + re.False(caps.SupportsPhase(2)) + re.False(caps.SupportsPhase(3)) + + caps2 := NewFullEncodingCapabilities(2) // Version 3 + re.True(caps2.SupportsPhase(1)) + re.True(caps2.SupportsPhase(2)) + re.True(caps2.SupportsPhase(3)) +} + +func TestClusterSupportsPhase(t *testing.T) { + re := require.New(t) + + stores := []*StoreEncodingCapabilities{ + NewFullEncodingCapabilities(1), + NewFullEncodingCapabilities(2), + } + status := ComputeClusterStatus(stores, 2) + re.True(status.ClusterSupportsPhase(1)) + re.True(status.ClusterSupportsPhase(2)) + re.True(status.ClusterSupportsPhase(3)) + + // Mixed: one v1, one v3 → min is 1 + stores2 := []*StoreEncodingCapabilities{ + NewStoreEncodingCapabilities(1), // v1 + NewFullEncodingCapabilities(2), // v3 + } + status2 := ComputeClusterStatus(stores2, 2) + re.True(status2.ClusterSupportsPhase(1)) + re.False(status2.ClusterSupportsPhase(2)) + re.False(status2.ClusterSupportsPhase(3)) +} + +func TestCheckOperationFeasibility(t *testing.T) { + re := require.New(t) + + stores := []*StoreEncodingCapabilities{ + NewFullEncodingCapabilities(1), + NewFullEncodingCapabilities(2), + NewFullEncodingCapabilities(3), + } + + result := CheckOperationFeasibility(stores, 3, CapEncodedStarJoin) + re.True(result.Feasible) + re.Equal(3, result.StoresReady) + re.Equal(3, result.StoresTotal) +} + +func TestCheckOperationFeasibilityPartial(t *testing.T) { + re := require.New(t) + + stores := []*StoreEncodingCapabilities{ + NewFullEncodingCapabilities(1), // has star join + NewStoreEncodingCapabilities(2), // no star join + } + + result := CheckOperationFeasibility(stores, 2, CapEncodedStarJoin) + re.False(result.Feasible) + re.Equal(1, result.StoresReady) + re.Equal(2, result.StoresTotal) +} + +func TestAllOperationsFeasibility(t *testing.T) { + re := require.New(t) + + stores := []*StoreEncodingCapabilities{ + NewFullEncodingCapabilities(1), + NewFullEncodingCapabilities(2), + } + + results := AllOperationsFeasibility(stores, 2) + re.Len(results, 5) // 5 operations + for _, r := range results { + re.True(r.Feasible) + re.Equal(2, r.StoresReady) + } +} + +func TestPhaseCapabilities(t *testing.T) { + re := require.New(t) + + re.Len(PhaseCapabilities[1], 1) + re.Len(PhaseCapabilities[2], 3) + re.Len(PhaseCapabilities[3], 5) + re.Contains(PhaseCapabilities[3], CapEncodedStarJoin) +}