From c7b64fe706eca5aaeb8c28503ed7f07b77adf082 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Mon, 23 Mar 2026 10:57:46 +0800 Subject: [PATCH 01/20] feat(eventstore): support encrypted persisted events Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 106 ++++++++++++++++--- logservice/eventstore/event_store_test.go | 118 ++++++++++++++++++++++ 2 files changed, 209 insertions(+), 15 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 82954a1cb9..153a6df52e 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -34,6 +34,7 @@ import ( "github.com/pingcap/ticdc/pkg/common" appcontext "github.com/pingcap/ticdc/pkg/common/context" "github.com/pingcap/ticdc/pkg/config" + "github.com/pingcap/ticdc/pkg/encryption" "github.com/pingcap/ticdc/pkg/messaging" "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/node" @@ -116,6 +117,8 @@ type dispatcherStat struct { resolvedTs atomic.Uint64 // the max ts of events which is not needed by this dispatcher checkpointTs uint64 + // keyspaceID for encryption (0 means default/classic) + keyspaceID uint32 // the difference between `subStat`, `pendingSubStat` and `removingSubStat`: // 1) if there is no existing subscriptions which can be reused, // or there is a existing subscription with exact span match, @@ -182,9 +185,10 @@ type subscriptionStat struct { type subscriptionStats map[logpuller.SubscriptionID]*subscriptionStat type eventWithCallback struct { - subID logpuller.SubscriptionID - tableID int64 - kvs []common.RawKVEntry + subID logpuller.SubscriptionID + tableID int64 + keyspaceID uint32 + kvs []common.RawKVEntry // kv with commitTs <= currentResolvedTs will be filtered out currentResolvedTs uint64 enqueueTimeNano int64 @@ -238,6 +242,8 @@ type eventStore struct { compressionThreshold int // enableZstdCompression controls whether to enable zstd compression for large values. enableZstdCompression bool + // encryptionManager for encrypting/decrypting data (optional). + encryptionManager encryption.EncryptionManager } const ( @@ -258,6 +264,9 @@ func New( log.Panic("fail to remove path", zap.String("path", dbPath), zap.Error(err)) } + // Try to get encryption manager from appcontext (optional) + encMgr, _ := appcontext.TryGetService[encryption.EncryptionManager](appcontext.EncryptionManager) + store := &eventStore{ pdClock: appcontext.GetService[pdutil.Clock](appcontext.DefaultPDClock), subClient: subClient, @@ -280,6 +289,7 @@ func New( }, compressionThreshold: config.GetGlobalServerConfig().Debug.EventStore.CompressionThreshold, enableZstdCompression: config.GetGlobalServerConfig().Debug.EventStore.EnableZstdCompression, + encryptionManager: encMgr, } store.gcManager = newGCManager(store.dbs, deleteDataRange, compactDataRange) @@ -484,6 +494,7 @@ func (e *eventStore) RegisterDispatcher( dispatcherID: dispatcherID, tableSpan: dispatcherSpan, checkpointTs: startTs, + keyspaceID: dispatcherSpan.KeyspaceID, } stat.resolvedTs.Store(startTs) @@ -610,6 +621,7 @@ func (e *eventStore) RegisterDispatcher( subStat.eventCh.Push(eventWithCallback{ subID: subStat.subID, tableID: subStat.tableSpan.TableID, + keyspaceID: subStat.tableSpan.KeyspaceID, kvs: kvs, currentResolvedTs: subStat.resolvedTs.Load(), enqueueTimeNano: now.UnixNano(), @@ -901,16 +913,18 @@ func (e *eventStore) GetIterator(dispatcherID common.DispatcherID, dataRange com } return &eventStoreIter{ - tableSpan: stat.tableSpan, - needCheckSpan: needCheckSpan, - innerIter: iter, - prevStartTs: 0, - prevCommitTs: 0, - startTs: dataRange.CommitTsStart, - endTs: dataRange.CommitTsEnd, - rowCount: 0, - decoder: decoder, - decoderPool: e.decoderPool, + tableSpan: stat.tableSpan, + needCheckSpan: needCheckSpan, + innerIter: iter, + prevStartTs: 0, + prevCommitTs: 0, + startTs: dataRange.CommitTsStart, + endTs: dataRange.CommitTsEnd, + rowCount: 0, + decoder: decoder, + decoderPool: e.decoderPool, + encryptionManager: e.encryptionManager, + keyspaceID: stat.keyspaceID, } } @@ -1319,7 +1333,58 @@ func (e *eventStore) writeEvents( valueBytesAfter := valueBytesBefore keyLen := encodedKeyLen(kv) - if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { + if e.encryptionManager != nil { + if cap(rawBuf) < int(valueBytesBefore) { + rawBuf = make([]byte, 0, int(valueBytesBefore)) + } else { + rawBuf = rawBuf[:0] + } + rawValue := kv.EncodeTo(rawBuf) + value := rawValue + if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { + maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) + if cap(dstBuf) < maxEncodedSize { + dstBuf = make([]byte, 0, maxEncodedSize) + } else { + dstBuf = dstBuf[:0] + } + value = encoder.EncodeAll(rawValue, dstBuf) + valueBytesAfter = int64(len(value)) + compressionType = CompressionZSTD + metrics.EventStoreCompressedRowsCount.Inc() + } + + // Encrypt if encryption is enabled (after compression) + encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) + if err != nil { + log.Error("encrypt event value failed", + zap.Uint32("keyspaceID", event.keyspaceID), + zap.Uint64("subID", uint64(event.subID)), + zap.Int64("tableID", event.tableID), + zap.Error(err)) + return err + } + + op := batch.SetDeferred(keyLen, len(encryptedValue)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + copiedValueLen := copy(op.Value, encryptedValue) + op.Value = op.Value[:copiedValueLen] + if copiedValueLen != len(encryptedValue) { + return fmt.Errorf("encrypted raw kv entry size mismatch, expected %d, got %d", + len(encryptedValue), copiedValueLen) + } + if err := op.Finish(); err != nil { + return err + } + rawBuf = rawValue[:0] + if compressionType == CompressionZSTD { + dstBuf = value[:0] + } + } else if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { if cap(rawBuf) < int(valueBytesBefore) { rawBuf = make([]byte, 0, int(valueBytesBefore)) } else { @@ -1376,7 +1441,6 @@ func (e *eventStore) writeEvents( return err } } - totalValueBytesBefore += valueBytesBefore totalValueBytesAfter += valueBytesAfter } @@ -1418,6 +1482,9 @@ type eventStoreIter struct { decoder *zstd.Decoder decoderPool *sync.Pool decodeBuf []byte + // encryptionManager for decrypting data (optional, can be nil). + encryptionManager encryption.EncryptionManager + keyspaceID uint32 } func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { @@ -1429,6 +1496,15 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { key := iter.innerIter.Key() value := iter.innerIter.Value() + // Decrypt if encrypted data is detected and encryption manager is available + if encryption.IsEncrypted(value) && iter.encryptionManager != nil { + decryptedValue, err := iter.encryptionManager.DecryptData(context.Background(), iter.keyspaceID, value) + if err != nil { + log.Panic("failed to decrypt value", zap.Error(err)) + } + value = decryptedValue + } + _, compressionType := DecodeKeyMetas(key) var decodedValue []byte if compressionType == CompressionZSTD { diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index f2cddd8ff8..4d52c42137 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -16,6 +16,7 @@ package eventstore import ( "bytes" "context" + "errors" "fmt" "os" "sync" @@ -30,6 +31,7 @@ import ( "github.com/pingcap/ticdc/pkg/common" appcontext "github.com/pingcap/ticdc/pkg/common/context" "github.com/pingcap/ticdc/pkg/config" + "github.com/pingcap/ticdc/pkg/encryption" "github.com/pingcap/ticdc/pkg/messaging" "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/pdutil" @@ -48,6 +50,31 @@ type mockSubscriptionClient struct { subscriptions map[logpuller.SubscriptionID]*mockSubscriptionStat } +type spyEncryptionManager struct { + encryptKeyspaceID uint32 + decryptKeyspaceID uint32 + encryptCalls int + decryptCalls int +} + +func (m *spyEncryptionManager) EncryptData(ctx context.Context, keyspaceID uint32, data []byte) ([]byte, error) { + m.encryptKeyspaceID = keyspaceID + m.encryptCalls++ + encrypted := make([]byte, encryption.EncryptionHeaderSize+len(data)) + encrypted[0] = 0x01 + copy(encrypted[encryption.EncryptionHeaderSize:], data) + return encrypted, nil +} + +func (m *spyEncryptionManager) DecryptData(ctx context.Context, keyspaceID uint32, encryptedData []byte) ([]byte, error) { + m.decryptKeyspaceID = keyspaceID + m.decryptCalls++ + if len(encryptedData) < encryption.EncryptionHeaderSize { + return nil, errors.New("encrypted data too short") + } + return encryptedData[encryption.EncryptionHeaderSize:], nil +} + func NewMockSubscriptionClient() logpuller.SubscriptionClient { return &mockSubscriptionClient{ subscriptions: make(map[logpuller.SubscriptionID]*mockSubscriptionStat), @@ -181,6 +208,97 @@ func TestEventStoreInteractionWithSubClient(t *testing.T) { } } +func TestEventStoreUsesKeyspaceIDForEncryption(t *testing.T) { + subClient, store := newEventStoreForTest(fmt.Sprintf("/tmp/%s", t.Name())) + es := store.(*eventStore) + defer es.Close(context.Background()) + + spy := &spyEncryptionManager{} + es.encryptionManager = spy + + dispatcherID := common.NewDispatcherID() + cfID := common.NewChangefeedID4Test("default", "test-cf") + span := &heartbeatpb.TableSpan{ + TableID: 1, + StartKey: []byte("a"), + EndKey: []byte("z"), + KeyspaceID: 42, + } + ok := store.RegisterDispatcher(cfID, dispatcherID, span, 0, func(uint64, uint64) {}, false, false) + require.True(t, ok) + + es.dispatcherMeta.RLock() + stat := es.dispatcherMeta.dispatcherStats[dispatcherID] + subStat := stat.subStat + if subStat == nil { + subStat = stat.pendingSubStat + } + es.dispatcherMeta.RUnlock() + require.NotNil(t, subStat) + + smallKV := common.RawKVEntry{ + OpType: common.OpTypePut, + CRTs: 10, + StartTs: 5, + Key: []byte("k-small"), + Value: []byte("v"), + } + largeKV := common.RawKVEntry{ + OpType: common.OpTypePut, + CRTs: 11, + StartTs: 6, + Key: []byte("k-large"), + Value: bytes.Repeat([]byte("v"), es.compressionThreshold+1), + } + encoder, err := zstd.NewWriter(nil) + require.NoError(t, err) + defer encoder.Close() + + events := []eventWithCallback{ + { + subID: subStat.subID, + tableID: subStat.tableSpan.TableID, + keyspaceID: subStat.tableSpan.KeyspaceID, + kvs: []common.RawKVEntry{smallKV, largeKV}, + currentResolvedTs: 0, + callback: func() {}, + }, + } + var compressionBuf []byte + var rawValueBuf []byte + err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + require.NoError(t, err) + require.Equal(t, uint32(42), spy.encryptKeyspaceID) + require.Equal(t, 2, spy.encryptCalls) + + subStat.resolvedTs.Store(largeKV.CRTs) + dataRange := common.DataRange{ + Span: span, + CommitTsStart: 0, + CommitTsEnd: largeKV.CRTs, + } + iter := es.GetIterator(dispatcherID, dataRange) + require.NotNil(t, iter) + + readValues := make(map[string][]byte) + for { + entry, ok := iter.Next() + if !ok { + break + } + readValues[string(entry.Key)] = entry.Value + } + require.Len(t, readValues, 2) + require.Equal(t, smallKV.Value, readValues[string(smallKV.Key)]) + require.Equal(t, largeKV.Value, readValues[string(largeKV.Key)]) + require.Equal(t, uint32(42), spy.decryptKeyspaceID) + require.Equal(t, 2, spy.decryptCalls) + + _, err = iter.Close() + require.NoError(t, err) + subClient.(*mockSubscriptionClient).Unsubscribe(subStat.subID) +} + func markSubStatsInitializedForTest(store EventStore, tableID int64) { es := store.(*eventStore) subStats := es.dispatcherMeta.tableStats[tableID] From 9aafb69de5be03ff1cac787ed516f0473e406d86 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Wed, 29 Apr 2026 09:37:52 +0800 Subject: [PATCH 02/20] feat(eventstore): remove encryption error logging - Remove detailed error logging when encryption fails in writeEvents - Encryption errors are now handled by returning the error directly - Reduces log noise while maintaining error propagation Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 5 ----- 1 file changed, 5 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 153a6df52e..5c040e37d1 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1357,11 +1357,6 @@ func (e *eventStore) writeEvents( // Encrypt if encryption is enabled (after compression) encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) if err != nil { - log.Error("encrypt event value failed", - zap.Uint32("keyspaceID", event.keyspaceID), - zap.Uint64("subID", uint64(event.subID)), - zap.Int64("tableID", event.tableID), - zap.Error(err)) return err } From f55924a81f685cff4fa96962d13f291d9c66c6dc Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Wed, 6 May 2026 10:57:25 +0800 Subject: [PATCH 03/20] eventstore: reorder write branches for fast path Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 48 ++++++++++++++-------------- 1 file changed, 24 insertions(+), 24 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 5c040e37d1..495accadf9 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1328,12 +1328,30 @@ func (e *eventStore) writeEvents( continue } - compressionType := CompressionNone valueBytesBefore := kv.GetSize() valueBytesAfter := valueBytesBefore keyLen := encodedKeyLen(kv) - - if e.encryptionManager != nil { + compressionType := CompressionNone + needCompress := e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) + if e.encryptionManager == nil && !needCompress { + // SetDeferred reserves the final Pebble batch space up front, so + // the uncompressed raw KV can be encoded directly into op.Value + // without allocating a separate value slice and copying it later. + op := batch.SetDeferred(keyLen, int(valueBytesBefore)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, CompressionNone) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + op.Value = kv.EncodeTo(op.Value[:0]) + if len(op.Value) != int(valueBytesBefore) { + return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", + valueBytesBefore, len(op.Value)) + } + if err := op.Finish(); err != nil { + return err + } + } else if e.encryptionManager != nil { if cap(rawBuf) < int(valueBytesBefore) { rawBuf = make([]byte, 0, int(valueBytesBefore)) } else { @@ -1341,7 +1359,7 @@ func (e *eventStore) writeEvents( } rawValue := kv.EncodeTo(rawBuf) value := rawValue - if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { + if needCompress { maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) if cap(dstBuf) < maxEncodedSize { dstBuf = make([]byte, 0, maxEncodedSize) @@ -1376,10 +1394,10 @@ func (e *eventStore) writeEvents( return err } rawBuf = rawValue[:0] - if compressionType == CompressionZSTD { + if needCompress { dstBuf = value[:0] } - } else if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { + } else { if cap(rawBuf) < int(valueBytesBefore) { rawBuf = make([]byte, 0, int(valueBytesBefore)) } else { @@ -1417,24 +1435,6 @@ func (e *eventStore) writeEvents( } rawBuf = rawValue[:0] dstBuf = value[:0] - } else { - // SetDeferred reserves the final Pebble batch space up front, so - // the uncompressed raw KV can be encoded directly into op.Value - // without allocating a separate value slice and copying it later. - op := batch.SetDeferred(keyLen, int(valueBytesBefore)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - op.Value = kv.EncodeTo(op.Value[:0]) - if len(op.Value) != int(valueBytesBefore) { - return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", - valueBytesBefore, len(op.Value)) - } - if err := op.Finish(); err != nil { - return err - } } totalValueBytesBefore += valueBytesBefore totalValueBytesAfter += valueBytesAfter From c1020d04deaa820b6a1af130f85ae5c1980bd1eb Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Wed, 6 May 2026 11:09:13 +0800 Subject: [PATCH 04/20] eventstore: refactor write value preparation Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 231 ++++++++++++---------- logservice/eventstore/event_store_test.go | 108 ++++++++++ 2 files changed, 238 insertions(+), 101 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 495accadf9..1052d38de2 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1293,6 +1293,124 @@ func (e *eventStore) collectAndReportStoreMetrics() { metrics.EventStoreResolvedTsLagGauge.Set(eventStoreResolvedTsLagInSec) } +type preparedValueForWrite struct { + value []byte + compressionType CompressionType + valueBytesAfter int64 + rawBuf []byte + compressionBuf []byte + directWrite bool +} + +func (e *eventStore) prepareValueForWrite( + encoder *zstd.Encoder, + keyspaceID uint32, + kv *common.RawKVEntry, + rawBuf []byte, + compressionBuf []byte, +) (preparedValueForWrite, error) { + valueBytesBefore := kv.GetSize() + needCompress := e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) + if e.encryptionManager == nil && !needCompress { + return preparedValueForWrite{ + compressionType: CompressionNone, + valueBytesAfter: valueBytesBefore, + rawBuf: rawBuf, + compressionBuf: compressionBuf, + directWrite: true, + }, nil + } + + if cap(rawBuf) < int(valueBytesBefore) { + rawBuf = make([]byte, 0, int(valueBytesBefore)) + } else { + rawBuf = rawBuf[:0] + } + rawValue := kv.EncodeTo(rawBuf) + value := rawValue + prepared := preparedValueForWrite{ + compressionType: CompressionNone, + valueBytesAfter: valueBytesBefore, + rawBuf: rawValue[:0], + compressionBuf: compressionBuf, + } + + if needCompress { + maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) + if cap(compressionBuf) < maxEncodedSize { + compressionBuf = make([]byte, 0, maxEncodedSize) + } else { + compressionBuf = compressionBuf[:0] + } + compressedValue := encoder.EncodeAll(rawValue, compressionBuf) + value = compressedValue + prepared.compressionType = CompressionZSTD + prepared.valueBytesAfter = int64(len(compressedValue)) + prepared.compressionBuf = compressedValue[:0] + metrics.EventStoreCompressedRowsCount.Inc() + } + + if e.encryptionManager != nil { + encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), keyspaceID, value) + if err != nil { + return preparedValueForWrite{}, err + } + value = encryptedValue + } + + prepared.value = value + return prepared, nil +} + +func writePreparedValueToBatch( + batch *pebble.Batch, + subID logpuller.SubscriptionID, + tableID int64, + kv *common.RawKVEntry, + keyLen int, + compressionType CompressionType, + value []byte, +) error { + op := batch.SetDeferred(keyLen, len(value)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(subID), tableID, kv, compressionType) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + copiedValueLen := copy(op.Value, value) + op.Value = op.Value[:copiedValueLen] + if copiedValueLen != len(value) { + return fmt.Errorf("prepared raw kv entry size mismatch, expected %d, got %d", + len(value), copiedValueLen) + } + return op.Finish() +} + +func writeRawValueToBatch( + batch *pebble.Batch, + subID logpuller.SubscriptionID, + tableID int64, + kv *common.RawKVEntry, + keyLen int, + valueBytesBefore int64, +) error { + // SetDeferred reserves the final Pebble batch space up front, so the + // uncompressed raw KV can be encoded directly into op.Value without + // allocating a separate value slice and copying it later. + op := batch.SetDeferred(keyLen, int(valueBytesBefore)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(subID), tableID, kv, CompressionNone) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + op.Value = kv.EncodeTo(op.Value[:0]) + if len(op.Value) != int(valueBytesBefore) { + return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", + valueBytesBefore, len(op.Value)) + } + return op.Finish() +} + func (e *eventStore) writeEvents( db *pebble.DB, events []eventWithCallback, @@ -1329,115 +1447,26 @@ func (e *eventStore) writeEvents( } valueBytesBefore := kv.GetSize() - valueBytesAfter := valueBytesBefore keyLen := encodedKeyLen(kv) - compressionType := CompressionNone - needCompress := e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) - if e.encryptionManager == nil && !needCompress { - // SetDeferred reserves the final Pebble batch space up front, so - // the uncompressed raw KV can be encoded directly into op.Value - // without allocating a separate value slice and copying it later. - op := batch.SetDeferred(keyLen, int(valueBytesBefore)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, CompressionNone) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - op.Value = kv.EncodeTo(op.Value[:0]) - if len(op.Value) != int(valueBytesBefore) { - return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", - valueBytesBefore, len(op.Value)) - } - if err := op.Finish(); err != nil { - return err - } - } else if e.encryptionManager != nil { - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] - } - rawValue := kv.EncodeTo(rawBuf) - value := rawValue - if needCompress { - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(dstBuf) < maxEncodedSize { - dstBuf = make([]byte, 0, maxEncodedSize) - } else { - dstBuf = dstBuf[:0] - } - value = encoder.EncodeAll(rawValue, dstBuf) - valueBytesAfter = int64(len(value)) - compressionType = CompressionZSTD - metrics.EventStoreCompressedRowsCount.Inc() - } - - // Encrypt if encryption is enabled (after compression) - encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) - if err != nil { - return err - } - - op := batch.SetDeferred(keyLen, len(encryptedValue)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - copiedValueLen := copy(op.Value, encryptedValue) - op.Value = op.Value[:copiedValueLen] - if copiedValueLen != len(encryptedValue) { - return fmt.Errorf("encrypted raw kv entry size mismatch, expected %d, got %d", - len(encryptedValue), copiedValueLen) - } - if err := op.Finish(); err != nil { + prepared, err := e.prepareValueForWrite(encoder, event.keyspaceID, kv, rawBuf, dstBuf) + if err != nil { + return err + } + rawBuf = prepared.rawBuf + dstBuf = prepared.compressionBuf + if prepared.directWrite { + if err := writeRawValueToBatch(batch, event.subID, event.tableID, kv, keyLen, valueBytesBefore); err != nil { return err } - rawBuf = rawValue[:0] - if needCompress { - dstBuf = value[:0] - } } else { - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] - } - rawValue := kv.EncodeTo(rawBuf) - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(dstBuf) < maxEncodedSize { - dstBuf = make([]byte, 0, maxEncodedSize) - } else { - dstBuf = dstBuf[:0] - } - value := encoder.EncodeAll(rawValue, dstBuf) - valueBytesAfter = int64(len(value)) - compressionType = CompressionZSTD - metrics.EventStoreCompressedRowsCount.Inc() - // SetDeferred is a write path optimization. Now that the compressed - // value length is known, reserve the exact key/value space in the - // Pebble batch, encode the key directly into op.Key, and copy the - // compressed value into op.Value without building a temporary key. - op := batch.SetDeferred(keyLen, len(value)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - copiedValueLen := copy(op.Value, value) - op.Value = op.Value[:copiedValueLen] - if copiedValueLen != len(value) { - return fmt.Errorf("compressed raw kv entry size mismatch, expected %d, got %d", - len(value), copiedValueLen) - } - if err := op.Finish(); err != nil { + if err := writePreparedValueToBatch( + batch, event.subID, event.tableID, kv, keyLen, prepared.compressionType, prepared.value, + ); err != nil { return err } - rawBuf = rawValue[:0] - dstBuf = value[:0] } totalValueBytesBefore += valueBytesBefore - totalValueBytesAfter += valueBytesAfter + totalValueBytesAfter += prepared.valueBytesAfter } } if compressionBuf != nil { diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index 4d52c42137..b440deb423 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -75,6 +75,114 @@ func (m *spyEncryptionManager) DecryptData(ctx context.Context, keyspaceID uint3 return encryptedData[encryption.EncryptionHeaderSize:], nil } +func TestPrepareValueForWrite(t *testing.T) { + encoder, err := zstd.NewWriter(nil) + require.NoError(t, err) + defer encoder.Close() + + decoder, err := zstd.NewReader(nil) + require.NoError(t, err) + defer decoder.Close() + + const compressionThreshold = 128 + testCases := []struct { + name string + value []byte + enableCompression bool + enableEncryption bool + wantDirectWrite bool + wantCompressionType CompressionType + }{ + { + name: "fast path", + value: []byte("small-value"), + enableCompression: true, + enableEncryption: false, + wantDirectWrite: true, + wantCompressionType: CompressionNone, + }, + { + name: "compression only", + value: bytes.Repeat([]byte("a"), compressionThreshold+64), + enableCompression: true, + enableEncryption: false, + wantDirectWrite: false, + wantCompressionType: CompressionZSTD, + }, + { + name: "encryption only", + value: []byte("small-value"), + enableCompression: false, + enableEncryption: true, + wantDirectWrite: false, + wantCompressionType: CompressionNone, + }, + { + name: "compression then encryption", + value: bytes.Repeat([]byte("b"), compressionThreshold+64), + enableCompression: true, + enableEncryption: true, + wantDirectWrite: false, + wantCompressionType: CompressionZSTD, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + kv := common.RawKVEntry{ + OpType: common.OpTypePut, + StartTs: 100, + CRTs: 110, + Key: []byte(tc.name), + Value: tc.value, + } + rawBufSeed := []byte("keep-raw") + compressionBufSeed := []byte("keep-compression") + + var spy *spyEncryptionManager + var encryptionManager encryption.EncryptionManager + if tc.enableEncryption { + spy = &spyEncryptionManager{} + encryptionManager = spy + } + + store := &eventStore{ + compressionThreshold: compressionThreshold, + enableZstdCompression: tc.enableCompression, + encryptionManager: encryptionManager, + } + + prepared, err := store.prepareValueForWrite( + encoder, 42, &kv, append([]byte(nil), rawBufSeed...), append([]byte(nil), compressionBufSeed...), + ) + require.NoError(t, err) + require.Equal(t, tc.wantDirectWrite, prepared.directWrite) + require.Equal(t, tc.wantCompressionType, prepared.compressionType) + + if tc.wantDirectWrite { + require.Nil(t, prepared.value) + require.Equal(t, rawBufSeed, prepared.rawBuf) + require.Equal(t, compressionBufSeed, prepared.compressionBuf) + require.Equal(t, kv.GetSize(), prepared.valueBytesAfter) + return + } + + transformed := prepared.value + if tc.enableEncryption { + require.Equal(t, 1, spy.encryptCalls) + transformed, err = spy.DecryptData(context.Background(), 42, transformed) + require.NoError(t, err) + } + if tc.wantCompressionType == CompressionZSTD { + transformed, err = decoder.DecodeAll(transformed, nil) + require.NoError(t, err) + } + require.Equal(t, kv.Encode(), transformed) + require.Len(t, prepared.rawBuf, 0) + }) + } +} + func NewMockSubscriptionClient() logpuller.SubscriptionClient { return &mockSubscriptionClient{ subscriptions: make(map[logpuller.SubscriptionID]*mockSubscriptionStat), From dcb834c45e9dc8b5df6f0f6019d21ae11f2c7f54 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Wed, 6 May 2026 18:06:30 +0800 Subject: [PATCH 05/20] Revert "eventstore: refactor write value preparation" This reverts commit c1020d04deaa820b6a1af130f85ae5c1980bd1eb. --- logservice/eventstore/event_store.go | 231 ++++++++++------------ logservice/eventstore/event_store_test.go | 108 ---------- 2 files changed, 101 insertions(+), 238 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 1052d38de2..495accadf9 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1293,124 +1293,6 @@ func (e *eventStore) collectAndReportStoreMetrics() { metrics.EventStoreResolvedTsLagGauge.Set(eventStoreResolvedTsLagInSec) } -type preparedValueForWrite struct { - value []byte - compressionType CompressionType - valueBytesAfter int64 - rawBuf []byte - compressionBuf []byte - directWrite bool -} - -func (e *eventStore) prepareValueForWrite( - encoder *zstd.Encoder, - keyspaceID uint32, - kv *common.RawKVEntry, - rawBuf []byte, - compressionBuf []byte, -) (preparedValueForWrite, error) { - valueBytesBefore := kv.GetSize() - needCompress := e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) - if e.encryptionManager == nil && !needCompress { - return preparedValueForWrite{ - compressionType: CompressionNone, - valueBytesAfter: valueBytesBefore, - rawBuf: rawBuf, - compressionBuf: compressionBuf, - directWrite: true, - }, nil - } - - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] - } - rawValue := kv.EncodeTo(rawBuf) - value := rawValue - prepared := preparedValueForWrite{ - compressionType: CompressionNone, - valueBytesAfter: valueBytesBefore, - rawBuf: rawValue[:0], - compressionBuf: compressionBuf, - } - - if needCompress { - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(compressionBuf) < maxEncodedSize { - compressionBuf = make([]byte, 0, maxEncodedSize) - } else { - compressionBuf = compressionBuf[:0] - } - compressedValue := encoder.EncodeAll(rawValue, compressionBuf) - value = compressedValue - prepared.compressionType = CompressionZSTD - prepared.valueBytesAfter = int64(len(compressedValue)) - prepared.compressionBuf = compressedValue[:0] - metrics.EventStoreCompressedRowsCount.Inc() - } - - if e.encryptionManager != nil { - encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), keyspaceID, value) - if err != nil { - return preparedValueForWrite{}, err - } - value = encryptedValue - } - - prepared.value = value - return prepared, nil -} - -func writePreparedValueToBatch( - batch *pebble.Batch, - subID logpuller.SubscriptionID, - tableID int64, - kv *common.RawKVEntry, - keyLen int, - compressionType CompressionType, - value []byte, -) error { - op := batch.SetDeferred(keyLen, len(value)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(subID), tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - copiedValueLen := copy(op.Value, value) - op.Value = op.Value[:copiedValueLen] - if copiedValueLen != len(value) { - return fmt.Errorf("prepared raw kv entry size mismatch, expected %d, got %d", - len(value), copiedValueLen) - } - return op.Finish() -} - -func writeRawValueToBatch( - batch *pebble.Batch, - subID logpuller.SubscriptionID, - tableID int64, - kv *common.RawKVEntry, - keyLen int, - valueBytesBefore int64, -) error { - // SetDeferred reserves the final Pebble batch space up front, so the - // uncompressed raw KV can be encoded directly into op.Value without - // allocating a separate value slice and copying it later. - op := batch.SetDeferred(keyLen, int(valueBytesBefore)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(subID), tableID, kv, CompressionNone) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) - } - op.Value = kv.EncodeTo(op.Value[:0]) - if len(op.Value) != int(valueBytesBefore) { - return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", - valueBytesBefore, len(op.Value)) - } - return op.Finish() -} - func (e *eventStore) writeEvents( db *pebble.DB, events []eventWithCallback, @@ -1447,26 +1329,115 @@ func (e *eventStore) writeEvents( } valueBytesBefore := kv.GetSize() + valueBytesAfter := valueBytesBefore keyLen := encodedKeyLen(kv) - prepared, err := e.prepareValueForWrite(encoder, event.keyspaceID, kv, rawBuf, dstBuf) - if err != nil { - return err - } - rawBuf = prepared.rawBuf - dstBuf = prepared.compressionBuf - if prepared.directWrite { - if err := writeRawValueToBatch(batch, event.subID, event.tableID, kv, keyLen, valueBytesBefore); err != nil { + compressionType := CompressionNone + needCompress := e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) + if e.encryptionManager == nil && !needCompress { + // SetDeferred reserves the final Pebble batch space up front, so + // the uncompressed raw KV can be encoded directly into op.Value + // without allocating a separate value slice and copying it later. + op := batch.SetDeferred(keyLen, int(valueBytesBefore)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, CompressionNone) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + op.Value = kv.EncodeTo(op.Value[:0]) + if len(op.Value) != int(valueBytesBefore) { + return fmt.Errorf("encoded raw kv entry size mismatch, expected %d, got %d", + valueBytesBefore, len(op.Value)) + } + if err := op.Finish(); err != nil { + return err + } + } else if e.encryptionManager != nil { + if cap(rawBuf) < int(valueBytesBefore) { + rawBuf = make([]byte, 0, int(valueBytesBefore)) + } else { + rawBuf = rawBuf[:0] + } + rawValue := kv.EncodeTo(rawBuf) + value := rawValue + if needCompress { + maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) + if cap(dstBuf) < maxEncodedSize { + dstBuf = make([]byte, 0, maxEncodedSize) + } else { + dstBuf = dstBuf[:0] + } + value = encoder.EncodeAll(rawValue, dstBuf) + valueBytesAfter = int64(len(value)) + compressionType = CompressionZSTD + metrics.EventStoreCompressedRowsCount.Inc() + } + + // Encrypt if encryption is enabled (after compression) + encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) + if err != nil { return err } + + op := batch.SetDeferred(keyLen, len(encryptedValue)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + copiedValueLen := copy(op.Value, encryptedValue) + op.Value = op.Value[:copiedValueLen] + if copiedValueLen != len(encryptedValue) { + return fmt.Errorf("encrypted raw kv entry size mismatch, expected %d, got %d", + len(encryptedValue), copiedValueLen) + } + if err := op.Finish(); err != nil { + return err + } + rawBuf = rawValue[:0] + if needCompress { + dstBuf = value[:0] + } } else { - if err := writePreparedValueToBatch( - batch, event.subID, event.tableID, kv, keyLen, prepared.compressionType, prepared.value, - ); err != nil { + if cap(rawBuf) < int(valueBytesBefore) { + rawBuf = make([]byte, 0, int(valueBytesBefore)) + } else { + rawBuf = rawBuf[:0] + } + rawValue := kv.EncodeTo(rawBuf) + maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) + if cap(dstBuf) < maxEncodedSize { + dstBuf = make([]byte, 0, maxEncodedSize) + } else { + dstBuf = dstBuf[:0] + } + value := encoder.EncodeAll(rawValue, dstBuf) + valueBytesAfter = int64(len(value)) + compressionType = CompressionZSTD + metrics.EventStoreCompressedRowsCount.Inc() + // SetDeferred is a write path optimization. Now that the compressed + // value length is known, reserve the exact key/value space in the + // Pebble batch, encode the key directly into op.Key, and copy the + // compressed value into op.Value without building a temporary key. + op := batch.SetDeferred(keyLen, len(value)) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + copiedValueLen := copy(op.Value, value) + op.Value = op.Value[:copiedValueLen] + if copiedValueLen != len(value) { + return fmt.Errorf("compressed raw kv entry size mismatch, expected %d, got %d", + len(value), copiedValueLen) + } + if err := op.Finish(); err != nil { return err } + rawBuf = rawValue[:0] + dstBuf = value[:0] } totalValueBytesBefore += valueBytesBefore - totalValueBytesAfter += prepared.valueBytesAfter + totalValueBytesAfter += valueBytesAfter } } if compressionBuf != nil { diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index b440deb423..4d52c42137 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -75,114 +75,6 @@ func (m *spyEncryptionManager) DecryptData(ctx context.Context, keyspaceID uint3 return encryptedData[encryption.EncryptionHeaderSize:], nil } -func TestPrepareValueForWrite(t *testing.T) { - encoder, err := zstd.NewWriter(nil) - require.NoError(t, err) - defer encoder.Close() - - decoder, err := zstd.NewReader(nil) - require.NoError(t, err) - defer decoder.Close() - - const compressionThreshold = 128 - testCases := []struct { - name string - value []byte - enableCompression bool - enableEncryption bool - wantDirectWrite bool - wantCompressionType CompressionType - }{ - { - name: "fast path", - value: []byte("small-value"), - enableCompression: true, - enableEncryption: false, - wantDirectWrite: true, - wantCompressionType: CompressionNone, - }, - { - name: "compression only", - value: bytes.Repeat([]byte("a"), compressionThreshold+64), - enableCompression: true, - enableEncryption: false, - wantDirectWrite: false, - wantCompressionType: CompressionZSTD, - }, - { - name: "encryption only", - value: []byte("small-value"), - enableCompression: false, - enableEncryption: true, - wantDirectWrite: false, - wantCompressionType: CompressionNone, - }, - { - name: "compression then encryption", - value: bytes.Repeat([]byte("b"), compressionThreshold+64), - enableCompression: true, - enableEncryption: true, - wantDirectWrite: false, - wantCompressionType: CompressionZSTD, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - kv := common.RawKVEntry{ - OpType: common.OpTypePut, - StartTs: 100, - CRTs: 110, - Key: []byte(tc.name), - Value: tc.value, - } - rawBufSeed := []byte("keep-raw") - compressionBufSeed := []byte("keep-compression") - - var spy *spyEncryptionManager - var encryptionManager encryption.EncryptionManager - if tc.enableEncryption { - spy = &spyEncryptionManager{} - encryptionManager = spy - } - - store := &eventStore{ - compressionThreshold: compressionThreshold, - enableZstdCompression: tc.enableCompression, - encryptionManager: encryptionManager, - } - - prepared, err := store.prepareValueForWrite( - encoder, 42, &kv, append([]byte(nil), rawBufSeed...), append([]byte(nil), compressionBufSeed...), - ) - require.NoError(t, err) - require.Equal(t, tc.wantDirectWrite, prepared.directWrite) - require.Equal(t, tc.wantCompressionType, prepared.compressionType) - - if tc.wantDirectWrite { - require.Nil(t, prepared.value) - require.Equal(t, rawBufSeed, prepared.rawBuf) - require.Equal(t, compressionBufSeed, prepared.compressionBuf) - require.Equal(t, kv.GetSize(), prepared.valueBytesAfter) - return - } - - transformed := prepared.value - if tc.enableEncryption { - require.Equal(t, 1, spy.encryptCalls) - transformed, err = spy.DecryptData(context.Background(), 42, transformed) - require.NoError(t, err) - } - if tc.wantCompressionType == CompressionZSTD { - transformed, err = decoder.DecodeAll(transformed, nil) - require.NoError(t, err) - } - require.Equal(t, kv.Encode(), transformed) - require.Len(t, prepared.rawBuf, 0) - }) - } -} - func NewMockSubscriptionClient() logpuller.SubscriptionClient { return &mockSubscriptionClient{ subscriptions: make(map[logpuller.SubscriptionID]*mockSubscriptionStat), From c5d957ad4e77f75860e2871c448b85c4bc60a5a5 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Wed, 6 May 2026 18:40:11 +0800 Subject: [PATCH 06/20] eventstore,encryption: frame plaintext encryption-layer values Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 9 +- logservice/eventstore/event_store_test.go | 107 ++++++++++++++++++++++ logservice/eventstore/format.go | 42 +++++++-- logservice/eventstore/pebble_test.go | 8 ++ pkg/encryption/encryption_manager.go | 21 ++++- pkg/encryption/encryption_manager_test.go | 12 ++- pkg/encryption/format.go | 26 ++++-- 7 files changed, 201 insertions(+), 24 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 495accadf9..c922ee177f 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1379,7 +1379,7 @@ func (e *eventStore) writeEvents( } op := batch.SetDeferred(keyLen, len(encryptedValue)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + op.Key = encodeKeyToWithEncryptionLayer(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) if len(op.Key) != keyLen { return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", keyLen, len(op.Key)) @@ -1491,8 +1491,11 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { key := iter.innerIter.Key() value := iter.innerIter.Value() - // Decrypt if encrypted data is detected and encryption manager is available - if encryption.IsEncrypted(value) && iter.encryptionManager != nil { + if KeyUsesEncryptionLayer(key) { + if iter.encryptionManager == nil { + log.Panic("encountered encryption-layer value but no encryption manager is configured", + zap.Uint32("keyspaceID", iter.keyspaceID)) + } decryptedValue, err := iter.encryptionManager.DecryptData(context.Background(), iter.keyspaceID, value) if err != nil { log.Panic("failed to decrypt value", zap.Error(err)) diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index 4d52c42137..d78c5e362e 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -32,6 +32,7 @@ import ( appcontext "github.com/pingcap/ticdc/pkg/common/context" "github.com/pingcap/ticdc/pkg/config" "github.com/pingcap/ticdc/pkg/encryption" + cerrors "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/messaging" "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/pdutil" @@ -57,6 +58,24 @@ type spyEncryptionManager struct { decryptCalls int } +type unencryptedMetaManager struct{} + +func (m *unencryptedMetaManager) IsEncryptionEnabled(ctx context.Context, keyspaceID uint32) bool { + return false +} + +func (m *unencryptedMetaManager) GetCurrentDataKey(ctx context.Context, keyspaceID uint32) ([]byte, string, byte, error) { + return nil, "", 0, nil +} + +func (m *unencryptedMetaManager) GetDataKey(ctx context.Context, keyspaceID uint32, dataKeyID string) ([]byte, error) { + return nil, cerrors.ErrDataKeyNotFound.GenWithStackByArgs("data key not found") +} + +func (m *unencryptedMetaManager) Start(ctx context.Context) error { return nil } + +func (m *unencryptedMetaManager) Stop() {} + func (m *spyEncryptionManager) EncryptData(ctx context.Context, keyspaceID uint32, data []byte) ([]byte, error) { m.encryptKeyspaceID = keyspaceID m.encryptCalls++ @@ -299,6 +318,94 @@ func TestEventStoreUsesKeyspaceIDForEncryption(t *testing.T) { subClient.(*mockSubscriptionClient).Unsubscribe(subStat.subID) } +func TestEventStoreHandlesUnencryptedValuesFromEncryptionLayer(t *testing.T) { + restoreCompression := setZstdCompressionForTest(t, true) + defer restoreCompression() + + subClient, store := newEventStoreForTest(t.TempDir()) + es := store.(*eventStore) + defer es.Close(context.Background()) + + es.encryptionManager = encryption.NewEncryptionManager(&unencryptedMetaManager{}) + + dispatcherID := common.NewDispatcherID() + cfID := common.NewChangefeedID4Test("default", "test-cf") + span := &heartbeatpb.TableSpan{ + TableID: 1, + StartKey: []byte("a"), + EndKey: []byte("z"), + KeyspaceID: 42, + } + ok := store.RegisterDispatcher(cfID, dispatcherID, span, 0, func(uint64, uint64) {}, false, false) + require.True(t, ok) + + es.dispatcherMeta.RLock() + stat := es.dispatcherMeta.dispatcherStats[dispatcherID] + subStat := stat.subStat + if subStat == nil { + subStat = stat.pendingSubStat + } + es.dispatcherMeta.RUnlock() + require.NotNil(t, subStat) + + smallKV := common.RawKVEntry{ + OpType: common.OpTypePut, + CRTs: 10, + StartTs: 5, + Key: []byte("k-small"), + Value: []byte("v"), + } + largeKV := common.RawKVEntry{ + OpType: common.OpTypePut, + CRTs: 11, + StartTs: 6, + Key: []byte("k-large"), + Value: bytes.Repeat([]byte("v"), es.compressionThreshold+1), + } + encoder, err := zstd.NewWriter(nil) + require.NoError(t, err) + defer encoder.Close() + + events := []eventWithCallback{ + { + subID: subStat.subID, + tableID: subStat.tableSpan.TableID, + keyspaceID: subStat.tableSpan.KeyspaceID, + kvs: []common.RawKVEntry{smallKV, largeKV}, + currentResolvedTs: 0, + callback: func() {}, + }, + } + var compressionBuf []byte + var rawValueBuf []byte + err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + require.NoError(t, err) + + subStat.resolvedTs.Store(largeKV.CRTs) + iter := es.GetIterator(dispatcherID, common.DataRange{ + Span: span, + CommitTsStart: 0, + CommitTsEnd: largeKV.CRTs, + }) + require.NotNil(t, iter) + + readValues := make(map[string][]byte) + for { + entry, ok := iter.Next() + if !ok { + break + } + readValues[string(entry.Key)] = entry.Value + } + require.Len(t, readValues, 2) + require.Equal(t, smallKV.Value, readValues[string(smallKV.Key)]) + require.Equal(t, largeKV.Value, readValues[string(largeKV.Key)]) + + _, err = iter.Close() + require.NoError(t, err) + subClient.(*mockSubscriptionClient).Unsubscribe(subStat.subID) +} + func markSubStatsInitializedForTest(store EventStore, tableID int64) { es := store.(*eventStore) subStats := es.dispatcherMeta.tableStats[tableID] diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index bb8cefe2b5..efed2b4ff7 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -40,9 +40,10 @@ const ( const ( // Bitmask for DML order and compression type. - dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. - compressionMask = 0x00FF // Compression type is stored in the low 8 bits. - dmlOrderShift = 8 + dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. + compressionMask = 0x000F // Compression type is stored in the low 4 bits. + encryptionLayerMask = 0x0010 // Value went through encryption layer and carries a 4-byte header. + dmlOrderShift = 8 ) // EncodeKeyPrefix encodes uniqueID, tableID, CRTs and StartTs. @@ -81,14 +82,13 @@ func encodedKeyLen(event *common.RawKVEntry) int { return 8 + 8 + 8 + 8 + 1 + 1 + len(event.Key) } -// EncodeKeyTo appends an encoded event-store key to buf. -// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. -func EncodeKeyTo( +func encodeKeyTo( buf []byte, uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType, + usesEncryptionLayer bool, ) []byte { if event == nil { log.Panic("rawkv must not be nil", zap.Any("event", event)) @@ -104,11 +104,36 @@ func EncodeKeyTo( // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) combinedOrder := uint16(compressionType) | (uint16(dmlOrder) << dmlOrderShift) + if usesEncryptionLayer { + combinedOrder |= encryptionLayerMask + } buf = binary.BigEndian.AppendUint16(buf, combinedOrder) // key return append(buf, event.Key...) } +// EncodeKeyTo appends an encoded event-store key to buf. +// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. +func EncodeKeyTo( + buf []byte, + uniqueID uint64, + tableID int64, + event *common.RawKVEntry, + compressionType CompressionType, +) []byte { + return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, false) +} + +func encodeKeyToWithEncryptionLayer( + buf []byte, + uniqueID uint64, + tableID int64, + event *common.RawKVEntry, + compressionType CompressionType, +) []byte { + return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, true) +} + // EncodeKey encodes a key according to event. func EncodeKey(uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType) []byte { return EncodeKeyTo(make([]byte, 0, encodedKeyLen(event)), uniqueID, tableID, event, compressionType) @@ -120,6 +145,11 @@ func DecodeKeyMetas(key []byte) (DMLOrder, CompressionType) { return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) } +func KeyUsesEncryptionLayer(key []byte) bool { + combinedOrder := binary.BigEndian.Uint16(key[32:34]) + return combinedOrder&encryptionLayerMask != 0 +} + // getDMLOrder returns the order of the dml types: delete Date: Wed, 6 May 2026 20:12:20 +0800 Subject: [PATCH 07/20] eventstore,encryption: preserve legacy framed value reads Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 31 +++++++++--- logservice/eventstore/event_store_test.go | 60 +++++++++++++++++++++++ logservice/eventstore/format.go | 42 +++------------- logservice/eventstore/pebble_test.go | 8 --- pkg/encryption/encryption_manager_test.go | 2 +- pkg/encryption/format.go | 17 ++++--- pkg/encryption/format_test.go | 10 ++-- 7 files changed, 105 insertions(+), 65 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index c922ee177f..7256c94cb8 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -58,6 +58,7 @@ var ( metricEventStoreFirstReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("first") metricEventStoreNextReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("next") metricEventStoreCloseReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("close") + legacyZSTDFrameMagic = []byte{0x28, 0xB5, 0x2F, 0xFD} ) const subscriptionIdleTTL = time.Minute @@ -1379,7 +1380,7 @@ func (e *eventStore) writeEvents( } op := batch.SetDeferred(keyLen, len(encryptedValue)) - op.Key = encodeKeyToWithEncryptionLayer(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) if len(op.Key) != keyLen { return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", keyLen, len(op.Key)) @@ -1490,12 +1491,13 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { } key := iter.innerIter.Key() value := iter.innerIter.Value() + _, compressionType := DecodeKeyMetas(key) - if KeyUsesEncryptionLayer(key) { - if iter.encryptionManager == nil { - log.Panic("encountered encryption-layer value but no encryption manager is configured", - zap.Uint32("keyspaceID", iter.keyspaceID)) - } + shouldDecrypt := iter.encryptionManager != nil + if shouldDecrypt && compressionType == CompressionZSTD && iter.isLegacyCompressedRawKV(value) { + shouldDecrypt = false + } + if shouldDecrypt { decryptedValue, err := iter.encryptionManager.DecryptData(context.Background(), iter.keyspaceID, value) if err != nil { log.Panic("failed to decrypt value", zap.Error(err)) @@ -1503,7 +1505,6 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { value = decryptedValue } - _, compressionType := DecodeKeyMetas(key) var decodedValue []byte if compressionType == CompressionZSTD { var err error @@ -1561,6 +1562,22 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { return rawKV, isNewTxn } +// Legacy compressed values were written without the 4-byte encryption-layer +// header, so they still begin with a plain ZSTD frame. +func (iter *eventStoreIter) isLegacyCompressedRawKV(value []byte) bool { + if len(value) < len(legacyZSTDFrameMagic) || !bytes.Equal(value[:len(legacyZSTDFrameMagic)], legacyZSTDFrameMagic) { + return false + } + + decodedValue, err := iter.decoder.DecodeAll(value, nil) + if err != nil { + return false + } + + probe := &common.RawKVEntry{} + return probe.Decode(decodedValue) == nil +} + func (iter *eventStoreIter) Close() (int64, error) { if iter.innerIter == nil { log.Info("event store close nil iter", diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index d78c5e362e..3be170810d 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -81,6 +81,7 @@ func (m *spyEncryptionManager) EncryptData(ctx context.Context, keyspaceID uint3 m.encryptCalls++ encrypted := make([]byte, encryption.EncryptionHeaderSize+len(data)) encrypted[0] = 0x01 + encrypted[3] = 0x01 copy(encrypted[encryption.EncryptionHeaderSize:], data) return encrypted, nil } @@ -1290,6 +1291,65 @@ func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { require.Equal(t, int64(len(expectedValues)), rowCount) } +func TestEventStoreIterReadsLegacyCompressedValuesWithEncryptionManager(t *testing.T) { + restoreCfg := setZstdCompressionForTest(t, true) + defer restoreCfg() + + dir := t.TempDir() + _, storeInt := newEventStoreForTest(dir) + store := storeInt.(*eventStore) + defer store.Close(context.Background()) + + entry := common.RawKVEntry{ + OpType: common.OpTypePut, + StartTs: 100, + CRTs: 200, + Key: []byte("legacy-compressed"), + Value: bytes.Repeat([]byte("value"), store.compressionThreshold), + } + events := []eventWithCallback{{ + subID: 1, + tableID: 1, + kvs: []common.RawKVEntry{entry}, + callback: func() {}, + }} + + encoder, err := zstd.NewWriter(nil) + require.NoError(t, err) + defer encoder.Close() + + var compressionBuf []byte + var rawValueBuf []byte + err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + require.NoError(t, err) + + innerIter, err := store.dbs[0].NewIter(&pebble.IterOptions{}) + require.NoError(t, err) + decoder := store.decoderPool.Get().(*zstd.Decoder) + iter := &eventStoreIter{ + tableSpan: &heartbeatpb.TableSpan{ + TableID: 1, + StartKey: []byte{}, + EndKey: []byte{0xFF}, + }, + innerIter: innerIter, + decoder: decoder, + decoderPool: store.decoderPool, + encryptionManager: encryption.NewEncryptionManager(&unencryptedMetaManager{}), + keyspaceID: 42, + } + require.True(t, iter.innerIter.First()) + + readEntry, ok := iter.Next() + require.True(t, ok) + require.Equal(t, entry.Key, readEntry.Key) + require.Equal(t, entry.Value, readEntry.Value) + + rowCount, err := iter.Close() + require.NoError(t, err) + require.Equal(t, int64(1), rowCount) +} + func TestEventStoreGetIteratorConcurrently(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index efed2b4ff7..bb8cefe2b5 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -40,10 +40,9 @@ const ( const ( // Bitmask for DML order and compression type. - dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. - compressionMask = 0x000F // Compression type is stored in the low 4 bits. - encryptionLayerMask = 0x0010 // Value went through encryption layer and carries a 4-byte header. - dmlOrderShift = 8 + dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. + compressionMask = 0x00FF // Compression type is stored in the low 8 bits. + dmlOrderShift = 8 ) // EncodeKeyPrefix encodes uniqueID, tableID, CRTs and StartTs. @@ -82,13 +81,14 @@ func encodedKeyLen(event *common.RawKVEntry) int { return 8 + 8 + 8 + 8 + 1 + 1 + len(event.Key) } -func encodeKeyTo( +// EncodeKeyTo appends an encoded event-store key to buf. +// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. +func EncodeKeyTo( buf []byte, uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType, - usesEncryptionLayer bool, ) []byte { if event == nil { log.Panic("rawkv must not be nil", zap.Any("event", event)) @@ -104,36 +104,11 @@ func encodeKeyTo( // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) combinedOrder := uint16(compressionType) | (uint16(dmlOrder) << dmlOrderShift) - if usesEncryptionLayer { - combinedOrder |= encryptionLayerMask - } buf = binary.BigEndian.AppendUint16(buf, combinedOrder) // key return append(buf, event.Key...) } -// EncodeKeyTo appends an encoded event-store key to buf. -// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. -func EncodeKeyTo( - buf []byte, - uniqueID uint64, - tableID int64, - event *common.RawKVEntry, - compressionType CompressionType, -) []byte { - return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, false) -} - -func encodeKeyToWithEncryptionLayer( - buf []byte, - uniqueID uint64, - tableID int64, - event *common.RawKVEntry, - compressionType CompressionType, -) []byte { - return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, true) -} - // EncodeKey encodes a key according to event. func EncodeKey(uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType) []byte { return EncodeKeyTo(make([]byte, 0, encodedKeyLen(event)), uniqueID, tableID, event, compressionType) @@ -145,11 +120,6 @@ func DecodeKeyMetas(key []byte) (DMLOrder, CompressionType) { return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) } -func KeyUsesEncryptionLayer(key []byte) bool { - combinedOrder := binary.BigEndian.Uint16(key[32:34]) - return combinedOrder&encryptionLayerMask != 0 -} - // getDMLOrder returns the order of the dml types: delete Date: Thu, 7 May 2026 11:21:09 +0800 Subject: [PATCH 08/20] eventstore: restore encryption-layer key metadata Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 31 +++++--------------- logservice/eventstore/format.go | 42 ++++++++++++++++++++++++---- logservice/eventstore/pebble_test.go | 8 ++++++ 3 files changed, 51 insertions(+), 30 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 7256c94cb8..c922ee177f 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -58,7 +58,6 @@ var ( metricEventStoreFirstReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("first") metricEventStoreNextReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("next") metricEventStoreCloseReadDurationHistogram = metrics.EventStoreReadDurationHistogram.WithLabelValues("close") - legacyZSTDFrameMagic = []byte{0x28, 0xB5, 0x2F, 0xFD} ) const subscriptionIdleTTL = time.Minute @@ -1380,7 +1379,7 @@ func (e *eventStore) writeEvents( } op := batch.SetDeferred(keyLen, len(encryptedValue)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) + op.Key = encodeKeyToWithEncryptionLayer(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) if len(op.Key) != keyLen { return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", keyLen, len(op.Key)) @@ -1491,13 +1490,12 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { } key := iter.innerIter.Key() value := iter.innerIter.Value() - _, compressionType := DecodeKeyMetas(key) - shouldDecrypt := iter.encryptionManager != nil - if shouldDecrypt && compressionType == CompressionZSTD && iter.isLegacyCompressedRawKV(value) { - shouldDecrypt = false - } - if shouldDecrypt { + if KeyUsesEncryptionLayer(key) { + if iter.encryptionManager == nil { + log.Panic("encountered encryption-layer value but no encryption manager is configured", + zap.Uint32("keyspaceID", iter.keyspaceID)) + } decryptedValue, err := iter.encryptionManager.DecryptData(context.Background(), iter.keyspaceID, value) if err != nil { log.Panic("failed to decrypt value", zap.Error(err)) @@ -1505,6 +1503,7 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { value = decryptedValue } + _, compressionType := DecodeKeyMetas(key) var decodedValue []byte if compressionType == CompressionZSTD { var err error @@ -1562,22 +1561,6 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { return rawKV, isNewTxn } -// Legacy compressed values were written without the 4-byte encryption-layer -// header, so they still begin with a plain ZSTD frame. -func (iter *eventStoreIter) isLegacyCompressedRawKV(value []byte) bool { - if len(value) < len(legacyZSTDFrameMagic) || !bytes.Equal(value[:len(legacyZSTDFrameMagic)], legacyZSTDFrameMagic) { - return false - } - - decodedValue, err := iter.decoder.DecodeAll(value, nil) - if err != nil { - return false - } - - probe := &common.RawKVEntry{} - return probe.Decode(decodedValue) == nil -} - func (iter *eventStoreIter) Close() (int64, error) { if iter.innerIter == nil { log.Info("event store close nil iter", diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index bb8cefe2b5..efed2b4ff7 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -40,9 +40,10 @@ const ( const ( // Bitmask for DML order and compression type. - dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. - compressionMask = 0x00FF // Compression type is stored in the low 8 bits. - dmlOrderShift = 8 + dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. + compressionMask = 0x000F // Compression type is stored in the low 4 bits. + encryptionLayerMask = 0x0010 // Value went through encryption layer and carries a 4-byte header. + dmlOrderShift = 8 ) // EncodeKeyPrefix encodes uniqueID, tableID, CRTs and StartTs. @@ -81,14 +82,13 @@ func encodedKeyLen(event *common.RawKVEntry) int { return 8 + 8 + 8 + 8 + 1 + 1 + len(event.Key) } -// EncodeKeyTo appends an encoded event-store key to buf. -// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. -func EncodeKeyTo( +func encodeKeyTo( buf []byte, uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType, + usesEncryptionLayer bool, ) []byte { if event == nil { log.Panic("rawkv must not be nil", zap.Any("event", event)) @@ -104,11 +104,36 @@ func EncodeKeyTo( // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) combinedOrder := uint16(compressionType) | (uint16(dmlOrder) << dmlOrderShift) + if usesEncryptionLayer { + combinedOrder |= encryptionLayerMask + } buf = binary.BigEndian.AppendUint16(buf, combinedOrder) // key return append(buf, event.Key...) } +// EncodeKeyTo appends an encoded event-store key to buf. +// Format: uniqueID, tableID, CRTs, startTs, delete/update/insert, Key. +func EncodeKeyTo( + buf []byte, + uniqueID uint64, + tableID int64, + event *common.RawKVEntry, + compressionType CompressionType, +) []byte { + return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, false) +} + +func encodeKeyToWithEncryptionLayer( + buf []byte, + uniqueID uint64, + tableID int64, + event *common.RawKVEntry, + compressionType CompressionType, +) []byte { + return encodeKeyTo(buf, uniqueID, tableID, event, compressionType, true) +} + // EncodeKey encodes a key according to event. func EncodeKey(uniqueID uint64, tableID int64, event *common.RawKVEntry, compressionType CompressionType) []byte { return EncodeKeyTo(make([]byte, 0, encodedKeyLen(event)), uniqueID, tableID, event, compressionType) @@ -120,6 +145,11 @@ func DecodeKeyMetas(key []byte) (DMLOrder, CompressionType) { return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) } +func KeyUsesEncryptionLayer(key []byte) bool { + combinedOrder := binary.BigEndian.Uint16(key[32:34]) + return combinedOrder&encryptionLayerMask != 0 +} + // getDMLOrder returns the order of the dml types: delete Date: Thu, 7 May 2026 13:44:19 +0800 Subject: [PATCH 09/20] eventstore: deduplicate writeEvents encoding Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 117 +++++++++++++++------------ 1 file changed, 65 insertions(+), 52 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 447a40e080..8f7d365349 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1371,10 +1371,8 @@ func (e *eventStore) writeEvents( // the uncompressed raw KV can be encoded directly into op.Value // without allocating a separate value slice and copying it later. op := batch.SetDeferred(keyLen, int(valueBytesBefore)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, CompressionNone) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) + if err := encodeDeferredEventKey(op, keyLen, uint64(event.subID), event.tableID, kv, CompressionNone, false); err != nil { + return err } op.Value = kv.EncodeTo(op.Value[:0]) if len(op.Value) != int(valueBytesBefore) { @@ -1385,25 +1383,9 @@ func (e *eventStore) writeEvents( return err } } else if e.encryptionManager != nil { - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] - } - rawValue := kv.EncodeTo(rawBuf) - value := rawValue - if needCompress { - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(dstBuf) < maxEncodedSize { - dstBuf = make([]byte, 0, maxEncodedSize) - } else { - dstBuf = dstBuf[:0] - } - value = encoder.EncodeAll(rawValue, dstBuf) - valueBytesAfter = int64(len(value)) - compressionType = CompressionZSTD - metrics.EventStoreCompressedRowsCount.Inc() - } + var value []byte + _, value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, needCompress) + valueBytesAfter = int64(len(value)) // Encrypt if encryption is enabled (after compression) encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) @@ -1412,10 +1394,8 @@ func (e *eventStore) writeEvents( } op := batch.SetDeferred(keyLen, len(encryptedValue)) - op.Key = encodeKeyToWithEncryptionLayer(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) + if err := encodeDeferredEventKey(op, keyLen, uint64(event.subID), event.tableID, kv, compressionType, true); err != nil { + return err } copiedValueLen := copy(op.Value, encryptedValue) op.Value = op.Value[:copiedValueLen] @@ -1426,36 +1406,17 @@ func (e *eventStore) writeEvents( if err := op.Finish(); err != nil { return err } - rawBuf = rawValue[:0] - if needCompress { - dstBuf = value[:0] - } } else { - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] - } - rawValue := kv.EncodeTo(rawBuf) - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(dstBuf) < maxEncodedSize { - dstBuf = make([]byte, 0, maxEncodedSize) - } else { - dstBuf = dstBuf[:0] - } - value := encoder.EncodeAll(rawValue, dstBuf) + var value []byte + _, value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, true) valueBytesAfter = int64(len(value)) - compressionType = CompressionZSTD - metrics.EventStoreCompressedRowsCount.Inc() // SetDeferred is a write path optimization. Now that the compressed // value length is known, reserve the exact key/value space in the // Pebble batch, encode the key directly into op.Key, and copy the // compressed value into op.Value without building a temporary key. op := batch.SetDeferred(keyLen, len(value)) - op.Key = EncodeKeyTo(op.Key[:0], uint64(event.subID), event.tableID, kv, compressionType) - if len(op.Key) != keyLen { - return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", - keyLen, len(op.Key)) + if err := encodeDeferredEventKey(op, keyLen, uint64(event.subID), event.tableID, kv, compressionType, false); err != nil { + return err } copiedValueLen := copy(op.Value, value) op.Value = op.Value[:copiedValueLen] @@ -1466,8 +1427,6 @@ func (e *eventStore) writeEvents( if err := op.Finish(); err != nil { return err } - rawBuf = rawValue[:0] - dstBuf = value[:0] } totalValueBytesBefore += valueBytesBefore totalValueBytesAfter += valueBytesAfter @@ -1493,6 +1452,60 @@ func (e *eventStore) writeEvents( return err } +func encodeAndMaybeCompressValue( + kv *common.RawKVEntry, + encoder *zstd.Encoder, + rawBuf []byte, + dstBuf []byte, + needCompress bool, +) (rawValue []byte, value []byte, compressionType CompressionType, nextRawBuf []byte, nextDstBuf []byte) { + rawValue = ensureValueBuffer(rawBuf, int(kv.GetSize())) + rawValue = kv.EncodeTo(rawValue) + value = rawValue + compressionType = CompressionNone + nextRawBuf = rawValue[:0] + nextDstBuf = dstBuf + if !needCompress { + return rawValue, value, compressionType, nextRawBuf, nextDstBuf + } + + maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) + value = ensureValueBuffer(dstBuf, maxEncodedSize) + value = encoder.EncodeAll(rawValue, value) + compressionType = CompressionZSTD + nextDstBuf = value[:0] + metrics.EventStoreCompressedRowsCount.Inc() + return rawValue, value, compressionType, nextRawBuf, nextDstBuf +} + +func ensureValueBuffer(buf []byte, minCap int) []byte { + if cap(buf) < minCap { + return make([]byte, 0, minCap) + } + return buf[:0] +} + +func encodeDeferredEventKey( + op *pebble.DeferredBatchOp, + keyLen int, + subID uint64, + tableID int64, + kv *common.RawKVEntry, + compressionType CompressionType, + usesEncryptionLayer bool, +) error { + if usesEncryptionLayer { + op.Key = encodeKeyToWithEncryptionLayer(op.Key[:0], subID, tableID, kv, compressionType) + } else { + op.Key = EncodeKeyTo(op.Key[:0], subID, tableID, kv, compressionType) + } + if len(op.Key) != keyLen { + return fmt.Errorf("encoded event store key size mismatch, expected %d, got %d", + keyLen, len(op.Key)) + } + return nil +} + type eventStoreIter struct { tableSpan *heartbeatpb.TableSpan // true when need check whether data from `innerIter` is in `tableSpan` From a2d8300438594c4fc1b987455527e89978536559 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 14:43:28 +0800 Subject: [PATCH 10/20] eventstore: widen encryption-layer key mask Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 14 +++++++------- logservice/eventstore/format_test.go | 9 +++++++-- 2 files changed, 14 insertions(+), 9 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index c979bd66b8..bd93d2da50 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -43,14 +43,14 @@ const ( encodedKeyTxnCommitTsStart = 2 * encodedKeyUint64Len encodedKeyTxnCommitTsEnd = encodedKeyTxnCommitTsStart + encodedKeyUint64Len encodedKeyAttributesOffset = 4 * encodedKeyUint64Len - encodedKeyAttributesEnd = encodedKeyAttributesOffset + 2 + encodedKeyAttributesEnd = encodedKeyAttributesOffset + 4 ) const ( // Bitmask for DML order and compression type. dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. - compressionMask = 0x000F // Compression type is stored in the low 4 bits. - encryptionLayerMask = 0x0010 // Value went through encryption layer and carries a 4-byte header. + compressionMask = 0x00FF // Compression type is stored in the low 8 bits. + encryptionLayerMask = 0x10000 dmlOrderShift = 8 ) @@ -100,11 +100,11 @@ func encodeKeyTo( buf = binary.BigEndian.AppendUint64(buf, event.StartTs) // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) - combinedOrder := uint16(compressionType) | (uint16(dmlOrder) << dmlOrderShift) + combinedOrder := uint32(compressionType) | (uint32(dmlOrder) << dmlOrderShift) if usesEncryptionLayer { combinedOrder |= encryptionLayerMask } - buf = binary.BigEndian.AppendUint16(buf, combinedOrder) + buf = binary.BigEndian.AppendUint32(buf, combinedOrder) // key return append(buf, event.Key...) } @@ -138,12 +138,12 @@ func EncodeKey(uniqueID uint64, tableID int64, event *common.RawKVEntry, compres // DecodeKeyAttributes decodes compression type and dml order from the key. func DecodeKeyAttributes(key []byte) (DMLOrder, CompressionType) { - combinedOrder := binary.BigEndian.Uint16(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) + combinedOrder := binary.BigEndian.Uint32(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) } func KeyUsesEncryptionLayer(key []byte) bool { - combinedOrder := binary.BigEndian.Uint16(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) + combinedOrder := binary.BigEndian.Uint32(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) return combinedOrder&encryptionLayerMask != 0 } diff --git a/logservice/eventstore/format_test.go b/logservice/eventstore/format_test.go index 4adf3a8b63..1ca46524b0 100644 --- a/logservice/eventstore/format_test.go +++ b/logservice/eventstore/format_test.go @@ -38,14 +38,14 @@ func TestEventStoreKeyFormatGolden(t *testing.T) { } key := EncodeKey(uniqueID, tableID, event, CompressionZSTD) - expectedKey := mustDecodeHex(t, "010203040506070811121314151617182122232425262728313233343536373803014142") + expectedKey := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738000003014142") require.Equal(t, expectedKey, key) require.Equal(t, len(expectedKey), encodedKeyLen(event)) require.Equal(t, 16, encodedKeyTxnCommitTsStart) require.Equal(t, 24, encodedKeyTxnCommitTsEnd) require.Equal(t, 32, encodedKeyAttributesOffset) - require.Equal(t, 34, encodedKeyAttributesEnd) + require.Equal(t, 36, encodedKeyAttributesEnd) require.Equal(t, expectedKey[:encodedKeyTxnCommitTsEnd], encodeTxnCommitTsBoundaryKey(uniqueID, tableID, txnCommitTs)) @@ -59,6 +59,11 @@ func TestEventStoreKeyFormatGolden(t *testing.T) { dmlOrder, compressionType := DecodeKeyAttributes(key) require.Equal(t, DMLOrderInsert, dmlOrder) require.Equal(t, CompressionZSTD, compressionType) + + keyWithEncryptionLayer := encodeKeyToWithEncryptionLayer(make([]byte, 0, encodedKeyLen(event)), uniqueID, tableID, event, CompressionZSTD) + expectedKeyWithEncryptionLayer := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738000103014142") + require.Equal(t, expectedKeyWithEncryptionLayer, keyWithEncryptionLayer) + require.True(t, KeyUsesEncryptionLayer(keyWithEncryptionLayer)) } func mustDecodeHex(t *testing.T, s string) []byte { From 33733ad9ad2c18fb86926547db3616de57bdeda5 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 15:59:36 +0800 Subject: [PATCH 11/20] eventstore: simplify key attribute encoding Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 41 ++++++++++++++-------------- logservice/eventstore/format_test.go | 9 ++++-- 2 files changed, 27 insertions(+), 23 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index bd93d2da50..e3d36b7d12 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -22,7 +22,7 @@ import ( "go.uber.org/zap" ) -type DMLOrder uint16 +type DMLOrder uint8 const ( // DML type order, used for sorting. @@ -31,7 +31,7 @@ const ( DMLOrderInsert ) -type CompressionType uint16 +type CompressionType uint8 const ( CompressionNone CompressionType = iota @@ -39,19 +39,19 @@ const ( ) const ( - encodedKeyUint64Len = 8 - encodedKeyTxnCommitTsStart = 2 * encodedKeyUint64Len - encodedKeyTxnCommitTsEnd = encodedKeyTxnCommitTsStart + encodedKeyUint64Len - encodedKeyAttributesOffset = 4 * encodedKeyUint64Len - encodedKeyAttributesEnd = encodedKeyAttributesOffset + 4 + encodedKeyUint64Len = 8 + encodedKeyTxnCommitTsStart = 2 * encodedKeyUint64Len + encodedKeyTxnCommitTsEnd = encodedKeyTxnCommitTsStart + encodedKeyUint64Len + encodedKeyAttributesOffset = 4 * encodedKeyUint64Len + encodedKeyDMLOrderOffset = encodedKeyAttributesOffset + encodedKeyCompressionOffset = encodedKeyDMLOrderOffset + 1 + encodedKeyMaskOffset = encodedKeyCompressionOffset + 1 + encodedKeyAttributesEnd = encodedKeyMaskOffset + encodedKeyUint64Len ) const ( - // Bitmask for DML order and compression type. - dmlOrderMask = 0xFF00 // DML order is stored in the high 8 bits for sorting. - compressionMask = 0x00FF // Compression type is stored in the low 8 bits. - encryptionLayerMask = 0x10000 - dmlOrderShift = 8 + // encodedKeyMaskEncryption indicates the value passed through the encryption layer. + encodedKeyMaskEncryption uint64 = 1 << iota ) // encodeTxnCommitTsBoundaryKey encodes the event-store key boundary up to txnCommitTs. @@ -75,7 +75,7 @@ func encodeTxnCommitTsBoundaryKeyTo(buf []byte, uniqueID uint64, tableID int64, } func encodedKeyLen(event *common.RawKVEntry) int { - // uniqueID, tableID, txnCommitTs, txnStartTs, Put/Delete, CompressionType, Key + // uniqueID, tableID, txnCommitTs, txnStartTs, DMLOrder, CompressionType, Mask, Key return encodedKeyAttributesEnd + len(event.Key) } @@ -100,11 +100,13 @@ func encodeKeyTo( buf = binary.BigEndian.AppendUint64(buf, event.StartTs) // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) - combinedOrder := uint32(compressionType) | (uint32(dmlOrder) << dmlOrderShift) + buf = append(buf, byte(dmlOrder)) + buf = append(buf, byte(compressionType)) + mask := uint64(0) if usesEncryptionLayer { - combinedOrder |= encryptionLayerMask + mask |= encodedKeyMaskEncryption } - buf = binary.BigEndian.AppendUint32(buf, combinedOrder) + buf = binary.BigEndian.AppendUint64(buf, mask) // key return append(buf, event.Key...) } @@ -138,13 +140,12 @@ func EncodeKey(uniqueID uint64, tableID int64, event *common.RawKVEntry, compres // DecodeKeyAttributes decodes compression type and dml order from the key. func DecodeKeyAttributes(key []byte) (DMLOrder, CompressionType) { - combinedOrder := binary.BigEndian.Uint32(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) - return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) + return DMLOrder(key[encodedKeyDMLOrderOffset]), CompressionType(key[encodedKeyCompressionOffset]) } func KeyUsesEncryptionLayer(key []byte) bool { - combinedOrder := binary.BigEndian.Uint32(key[encodedKeyAttributesOffset:encodedKeyAttributesEnd]) - return combinedOrder&encryptionLayerMask != 0 + mask := binary.BigEndian.Uint64(key[encodedKeyMaskOffset:encodedKeyAttributesEnd]) + return mask&encodedKeyMaskEncryption != 0 } // decodeTxnCommitTsFromEncodedKey decodes txnCommitTs from an event-store key boundary. diff --git a/logservice/eventstore/format_test.go b/logservice/eventstore/format_test.go index 1ca46524b0..cf7fc47572 100644 --- a/logservice/eventstore/format_test.go +++ b/logservice/eventstore/format_test.go @@ -38,14 +38,16 @@ func TestEventStoreKeyFormatGolden(t *testing.T) { } key := EncodeKey(uniqueID, tableID, event, CompressionZSTD) - expectedKey := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738000003014142") + expectedKey := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738030100000000000000004142") require.Equal(t, expectedKey, key) require.Equal(t, len(expectedKey), encodedKeyLen(event)) require.Equal(t, 16, encodedKeyTxnCommitTsStart) require.Equal(t, 24, encodedKeyTxnCommitTsEnd) require.Equal(t, 32, encodedKeyAttributesOffset) - require.Equal(t, 36, encodedKeyAttributesEnd) + require.Equal(t, 33, encodedKeyCompressionOffset) + require.Equal(t, 34, encodedKeyMaskOffset) + require.Equal(t, 42, encodedKeyAttributesEnd) require.Equal(t, expectedKey[:encodedKeyTxnCommitTsEnd], encodeTxnCommitTsBoundaryKey(uniqueID, tableID, txnCommitTs)) @@ -59,9 +61,10 @@ func TestEventStoreKeyFormatGolden(t *testing.T) { dmlOrder, compressionType := DecodeKeyAttributes(key) require.Equal(t, DMLOrderInsert, dmlOrder) require.Equal(t, CompressionZSTD, compressionType) + require.False(t, KeyUsesEncryptionLayer(key)) keyWithEncryptionLayer := encodeKeyToWithEncryptionLayer(make([]byte, 0, encodedKeyLen(event)), uniqueID, tableID, event, CompressionZSTD) - expectedKeyWithEncryptionLayer := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738000103014142") + expectedKeyWithEncryptionLayer := mustDecodeHex(t, "0102030405060708111213141516171821222324252627283132333435363738030100000000000000014142") require.Equal(t, expectedKeyWithEncryptionLayer, keyWithEncryptionLayer) require.True(t, KeyUsesEncryptionLayer(keyWithEncryptionLayer)) } From b776befb13f947a9af73621e732f7aae17464e9c Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 16:28:03 +0800 Subject: [PATCH 12/20] eventstore: clarify key field offsets Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 46 +++++++++++++++++++--------- logservice/eventstore/format_test.go | 7 +++-- logservice/eventstore/pebble_test.go | 2 +- 3 files changed, 37 insertions(+), 18 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index e3d36b7d12..30476c2b45 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -39,14 +39,30 @@ const ( ) const ( - encodedKeyUint64Len = 8 - encodedKeyTxnCommitTsStart = 2 * encodedKeyUint64Len - encodedKeyTxnCommitTsEnd = encodedKeyTxnCommitTsStart + encodedKeyUint64Len - encodedKeyAttributesOffset = 4 * encodedKeyUint64Len - encodedKeyDMLOrderOffset = encodedKeyAttributesOffset + // Encoded event-store key layout: + // + // byte offset + // 0 8 16 24 32 33 34 42 + // |-------------------|------------------|------------------|------------------|--|--|------------------| + // | uniqueID (8B) | tableID (8B) | txnCommitTs (8B) | txnStartTs (8B) |DO|CT| mask (8B) | key... + // |-------------------|------------------|------------------|------------------|--|--|------------------| + // + // DO: DMLOrder (1 byte) + // CT: CompressionType (1 byte) + // + // Mask bits: + // bit 0: value passed through the encryption layer + // bit 1+: reserved + encodedKeyUniqueIDOffset = 0 + encodedKeyTableIDOffset = encodedKeyUniqueIDOffset + 8 + encodedKeyTxnCommitTsOffset = encodedKeyTableIDOffset + 8 + encodedKeyTxnStartTsOffset = encodedKeyTxnCommitTsOffset + 8 + encodedKeyDMLOrderOffset = encodedKeyTxnStartTsOffset + 8 encodedKeyCompressionOffset = encodedKeyDMLOrderOffset + 1 encodedKeyMaskOffset = encodedKeyCompressionOffset + 1 - encodedKeyAttributesEnd = encodedKeyMaskOffset + encodedKeyUint64Len + + encodedKeyAttributesOffset = encodedKeyDMLOrderOffset + encodedKeyAttributesEnd = encodedKeyMaskOffset + 8 ) const ( @@ -56,7 +72,7 @@ const ( // encodeTxnCommitTsBoundaryKey encodes the event-store key boundary up to txnCommitTs. func encodeTxnCommitTsBoundaryKey(uniqueID uint64, tableID int64, txnCommitTs uint64) []byte { - buf := make([]byte, encodedKeyTxnCommitTsEnd) + buf := make([]byte, encodedKeyTxnStartTsOffset) encodeTxnCommitTsBoundaryKeyTo(buf, uniqueID, tableID, txnCommitTs) return buf } @@ -64,14 +80,14 @@ func encodeTxnCommitTsBoundaryKey(uniqueID uint64, tableID int64, txnCommitTs ui func encodeScanLowerBound(uniqueID uint64, tableID int64, txnCommitTs uint64, txnStartTs uint64) []byte { buf := make([]byte, encodedKeyAttributesOffset) encodeTxnCommitTsBoundaryKeyTo(buf, uniqueID, tableID, txnCommitTs) - binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsEnd:encodedKeyAttributesOffset], txnStartTs) + binary.BigEndian.PutUint64(buf[encodedKeyTxnStartTsOffset:encodedKeyDMLOrderOffset], txnStartTs) return buf } func encodeTxnCommitTsBoundaryKeyTo(buf []byte, uniqueID uint64, tableID int64, txnCommitTs uint64) { - binary.BigEndian.PutUint64(buf[:encodedKeyUint64Len], uniqueID) - binary.BigEndian.PutUint64(buf[encodedKeyUint64Len:encodedKeyTxnCommitTsStart], uint64(tableID)) - binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsStart:encodedKeyTxnCommitTsEnd], txnCommitTs) + binary.BigEndian.PutUint64(buf[encodedKeyUniqueIDOffset:encodedKeyTableIDOffset], uniqueID) + binary.BigEndian.PutUint64(buf[encodedKeyTableIDOffset:encodedKeyTxnCommitTsOffset], uint64(tableID)) + binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsOffset:encodedKeyTxnStartTsOffset], txnCommitTs) } func encodedKeyLen(event *common.RawKVEntry) int { @@ -112,7 +128,9 @@ func encodeKeyTo( } // EncodeKeyTo appends an encoded event-store key to buf. -// Format: uniqueID, tableID, txnCommitTs, txnStartTs, delete/update/insert, Key. +// +// | uniqueID | tableID | txnCommitTs | txnStartTs | dmlOrder | compressionType | mask | key | +// | 8B | 8B | 8B | 8B | 1B | 1B | 8B | ... | func EncodeKeyTo( buf []byte, uniqueID uint64, @@ -152,10 +170,10 @@ func KeyUsesEncryptionLayer(key []byte) bool { // It works for both full event keys and DeleteRange boundary keys because both // contain uniqueID, tableID, and txnCommitTs as the first three fields. func decodeTxnCommitTsFromEncodedKey(key []byte) (uint64, bool) { - if len(key) < encodedKeyTxnCommitTsEnd { + if len(key) < encodedKeyTxnStartTsOffset { return 0, false } - return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsStart:encodedKeyTxnCommitTsEnd]), true + return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsOffset:encodedKeyTxnStartTsOffset]), true } // getDMLOrder returns the order of the dml types: delete Date: Thu, 7 May 2026 16:34:25 +0800 Subject: [PATCH 13/20] fix(eventstore): Correct diagram alignment in format.go - Fix typo in comment diagram for event log format - Adjust spacing to maintain proper alignment in byte offset visualization Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index 30476c2b45..0a5cea1a94 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -44,7 +44,7 @@ const ( // byte offset // 0 8 16 24 32 33 34 42 // |-------------------|------------------|------------------|------------------|--|--|------------------| - // | uniqueID (8B) | tableID (8B) | txnCommitTs (8B) | txnStartTs (8B) |DO|CT| mask (8B) | key... + // | uniqueID (8B) | tableID (8B) | txnCommitTs (8B) | txnStartTs (8B) |DO|CT| mask (8B) | key... // |-------------------|------------------|------------------|------------------|--|--|------------------| // // DO: DMLOrder (1 byte) From 67be55702a9f8525f99c78be11221f223b2277ae Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 17:25:35 +0800 Subject: [PATCH 14/20] eventstore: clarify key field lengths Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 37 +++++++++++++++++----------- logservice/eventstore/format_test.go | 19 +++++++++++++- 2 files changed, 41 insertions(+), 15 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index 0a5cea1a94..9ddf61edc1 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -53,16 +53,25 @@ const ( // Mask bits: // bit 0: value passed through the encryption layer // bit 1+: reserved - encodedKeyUniqueIDOffset = 0 - encodedKeyTableIDOffset = encodedKeyUniqueIDOffset + 8 - encodedKeyTxnCommitTsOffset = encodedKeyTableIDOffset + 8 - encodedKeyTxnStartTsOffset = encodedKeyTxnCommitTsOffset + 8 - encodedKeyDMLOrderOffset = encodedKeyTxnStartTsOffset + 8 - encodedKeyCompressionOffset = encodedKeyDMLOrderOffset + 1 - encodedKeyMaskOffset = encodedKeyCompressionOffset + 1 + encodedKeyUniqueIDLen = 8 + encodedKeyTableIDLen = 8 + encodedKeyTxnCommitTsLen = 8 + encodedKeyTxnStartTsLen = 8 + encodedKeyDMLOrderLen = 1 + encodedKeyCompressionLen = 1 + encodedKeyMaskLen = 8 + encodedKeyUniqueIDOffset = 0 + encodedKeyTableIDOffset = encodedKeyUniqueIDOffset + encodedKeyUniqueIDLen + encodedKeyTxnCommitTsOffset = encodedKeyTableIDOffset + encodedKeyTableIDLen + encodedKeyTxnStartTsOffset = encodedKeyTxnCommitTsOffset + encodedKeyTxnCommitTsLen + encodedKeyDMLOrderOffset = encodedKeyTxnStartTsOffset + encodedKeyTxnStartTsLen + encodedKeyCompressionOffset = encodedKeyDMLOrderOffset + encodedKeyDMLOrderLen + encodedKeyMaskOffset = encodedKeyCompressionOffset + encodedKeyCompressionLen + + encodedKeyAttributesLen = encodedKeyDMLOrderLen + encodedKeyCompressionLen + encodedKeyMaskLen encodedKeyAttributesOffset = encodedKeyDMLOrderOffset - encodedKeyAttributesEnd = encodedKeyMaskOffset + 8 + encodedKeyAttributesEnd = encodedKeyAttributesOffset + encodedKeyAttributesLen ) const ( @@ -80,14 +89,14 @@ func encodeTxnCommitTsBoundaryKey(uniqueID uint64, tableID int64, txnCommitTs ui func encodeScanLowerBound(uniqueID uint64, tableID int64, txnCommitTs uint64, txnStartTs uint64) []byte { buf := make([]byte, encodedKeyAttributesOffset) encodeTxnCommitTsBoundaryKeyTo(buf, uniqueID, tableID, txnCommitTs) - binary.BigEndian.PutUint64(buf[encodedKeyTxnStartTsOffset:encodedKeyDMLOrderOffset], txnStartTs) + binary.BigEndian.PutUint64(buf[encodedKeyTxnStartTsOffset:encodedKeyTxnStartTsOffset+encodedKeyTxnStartTsLen], txnStartTs) return buf } func encodeTxnCommitTsBoundaryKeyTo(buf []byte, uniqueID uint64, tableID int64, txnCommitTs uint64) { - binary.BigEndian.PutUint64(buf[encodedKeyUniqueIDOffset:encodedKeyTableIDOffset], uniqueID) - binary.BigEndian.PutUint64(buf[encodedKeyTableIDOffset:encodedKeyTxnCommitTsOffset], uint64(tableID)) - binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsOffset:encodedKeyTxnStartTsOffset], txnCommitTs) + binary.BigEndian.PutUint64(buf[encodedKeyUniqueIDOffset:encodedKeyUniqueIDOffset+encodedKeyUniqueIDLen], uniqueID) + binary.BigEndian.PutUint64(buf[encodedKeyTableIDOffset:encodedKeyTableIDOffset+encodedKeyTableIDLen], uint64(tableID)) + binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsOffset:encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen], txnCommitTs) } func encodedKeyLen(event *common.RawKVEntry) int { @@ -162,7 +171,7 @@ func DecodeKeyAttributes(key []byte) (DMLOrder, CompressionType) { } func KeyUsesEncryptionLayer(key []byte) bool { - mask := binary.BigEndian.Uint64(key[encodedKeyMaskOffset:encodedKeyAttributesEnd]) + mask := binary.BigEndian.Uint64(key[encodedKeyMaskOffset : encodedKeyMaskOffset+encodedKeyMaskLen]) return mask&encodedKeyMaskEncryption != 0 } @@ -173,7 +182,7 @@ func decodeTxnCommitTsFromEncodedKey(key []byte) (uint64, bool) { if len(key) < encodedKeyTxnStartTsOffset { return 0, false } - return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsOffset:encodedKeyTxnStartTsOffset]), true + return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsOffset : encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen]), true } // getDMLOrder returns the order of the dml types: delete Date: Thu, 7 May 2026 17:30:07 +0800 Subject: [PATCH 15/20] eventstore: move key layout comment Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index 9ddf61edc1..7cec5b656d 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -104,6 +104,10 @@ func encodedKeyLen(event *common.RawKVEntry) int { return encodedKeyAttributesEnd + len(event.Key) } +// encodeKeyTo appends an encoded event-store key to buf. +// +// | uniqueID | tableID | txnCommitTs | txnStartTs | dmlOrder | compressionType | mask | key | +// | 8B | 8B | 8B | 8B | 1B | 1B | 8B | ... | func encodeKeyTo( buf []byte, uniqueID uint64, @@ -137,9 +141,6 @@ func encodeKeyTo( } // EncodeKeyTo appends an encoded event-store key to buf. -// -// | uniqueID | tableID | txnCommitTs | txnStartTs | dmlOrder | compressionType | mask | key | -// | 8B | 8B | 8B | 8B | 1B | 1B | 8B | ... | func EncodeKeyTo( buf []byte, uniqueID uint64, From 3f3e35960086a4f89e71b453734f91ebc8cb8c1a Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 18:00:06 +0800 Subject: [PATCH 16/20] encryption: tighten encryption layer decoding Signed-off-by: tenfyzhong --- pkg/encryption/encryption_manager.go | 64 +++++++++++------------ pkg/encryption/encryption_manager_test.go | 24 +++++++++ pkg/encryption/format.go | 26 ++++----- pkg/encryption/format_test.go | 30 ++++------- 4 files changed, 76 insertions(+), 68 deletions(-) diff --git a/pkg/encryption/encryption_manager.go b/pkg/encryption/encryption_manager.go index 7916413c55..53026ff879 100644 --- a/pkg/encryption/encryption_manager.go +++ b/pkg/encryption/encryption_manager.go @@ -25,13 +25,15 @@ import ( // EncryptionManager is the main interface for encryption/decryption operations type EncryptionManager interface { // EncryptData encrypts data for a keyspace - // Returns data wrapped with a 4-byte header. The payload is encrypted when a - // data key is available, or left plaintext with a version-0 header otherwise. + // Returns a value that has passed through the encryption layer and is wrapped + // with a 4-byte header. The payload is encrypted when a data key is available, + // or left plaintext with a version-0 header otherwise. EncryptData(ctx context.Context, keyspaceID uint32, data []byte) ([]byte, error) // DecryptData decrypts data for a keyspace - // Automatically unwraps plaintext headers and decrypts encrypted payloads. - DecryptData(ctx context.Context, keyspaceID uint32, encryptedData []byte) ([]byte, error) + // The caller must ensure the key marks this value as passing through the + // encryption layer, so the value must carry the 4-byte encryption header. + DecryptData(ctx context.Context, keyspaceID uint32, layerData []byte) ([]byte, error) } type encryptionManager struct { @@ -119,41 +121,39 @@ func (m *encryptionManager) EncryptData(ctx context.Context, keyspaceID uint32, return result, nil } -// DecryptData decrypts data for a keyspace -func (m *encryptionManager) DecryptData(ctx context.Context, keyspaceID uint32, encryptedData []byte) ([]byte, error) { - if HasUnencryptedHeader(encryptedData) { - plaintext, err := DecodeUnencryptedData(encryptedData) - if err != nil { - log.Warn("failed to decode unencrypted data header", - zap.Uint32("keyspaceID", keyspaceID), - zap.Int("encryptedSize", len(encryptedData)), - zap.Error(err)) - return nil, cerrors.ErrDecryptionFailed.Wrap(err) +func decodeEncryptionLayerData(layerData []byte) (byte, string, []byte, error) { + version, dataKeyID, payload, err := DecodeEncryptedData(layerData) + if err != nil { + return 0, "", nil, err + } + + if version == VersionUnencrypted { + if !dataKeyIDIsZero(layerData) { + return 0, "", nil, cerrors.ErrDecodeFailed.GenWithStackByArgs("invalid unencrypted data header") } - return plaintext, nil + return version, "", payload, nil } - // Check if data is encrypted - if !IsEncrypted(encryptedData) { - // Data is not encrypted, return as-is (backward compatibility) - log.Debug("data is not encrypted", - zap.Uint32("keyspaceID", keyspaceID)) - return encryptedData, nil + if dataKeyIDIsZero(layerData) { + return 0, "", nil, cerrors.ErrDecodeFailed.GenWithStackByArgs("invalid encrypted data header") } - // Decode encryption header - version, dataKeyID, dataWithIV, err := DecodeEncryptedData(encryptedData) + return version, dataKeyID, payload, nil +} + +// DecryptData decrypts data for a keyspace +func (m *encryptionManager) DecryptData(ctx context.Context, keyspaceID uint32, layerData []byte) ([]byte, error) { + version, dataKeyID, payload, err := decodeEncryptionLayerData(layerData) if err != nil { - log.Warn("failed to decode encrypted data header", + log.Warn("failed to decode encryption layer data", zap.Uint32("keyspaceID", keyspaceID), - zap.Int("encryptedSize", len(encryptedData)), + zap.Int("valueSize", len(layerData)), zap.Error(err)) return nil, cerrors.ErrDecryptionFailed.Wrap(err) } if version == VersionUnencrypted { - // Should not happen if IsEncrypted returned true, but handle it anyway - return dataWithIV, nil + return payload, nil } dataKey, err := m.metaManager.GetDataKey(ctx, keyspaceID, dataKeyID) @@ -178,18 +178,18 @@ func (m *encryptionManager) DecryptData(ctx context.Context, keyspaceID uint32, // Extract IV from the beginning of data ivSize := cipherImpl.IVSize() - if len(dataWithIV) < ivSize { + if len(payload) < ivSize { log.Warn("encrypted data too short for IV", zap.Uint32("keyspaceID", keyspaceID), zap.Uint8("version", version), zap.Binary("dataKeyID", []byte(dataKeyID)), - zap.Int("dataWithIVSize", len(dataWithIV)), + zap.Int("dataWithIVSize", len(payload)), zap.Int("expectedIVSize", ivSize)) return nil, cerrors.ErrDecryptionFailed.GenWithStackByArgs("data too short for IV") } - iv := dataWithIV[:ivSize] - encryptedDataOnly := dataWithIV[ivSize:] + iv := payload[:ivSize] + encryptedDataOnly := payload[ivSize:] // Decrypt data plaintext, err := cipherImpl.Decrypt(encryptedDataOnly, dataKey, iv) @@ -205,7 +205,7 @@ func (m *encryptionManager) DecryptData(ctx context.Context, keyspaceID uint32, log.Debug("data decrypted successfully", zap.Uint32("keyspaceID", keyspaceID), zap.String("dataKeyID", dataKeyID), - zap.Int("encryptedSize", len(encryptedData)), + zap.Int("encryptedSize", len(layerData)), zap.Int("plaintextSize", len(plaintext))) return plaintext, nil diff --git a/pkg/encryption/encryption_manager_test.go b/pkg/encryption/encryption_manager_test.go index d111f06680..c984688f5a 100644 --- a/pkg/encryption/encryption_manager_test.go +++ b/pkg/encryption/encryption_manager_test.go @@ -160,3 +160,27 @@ func TestEncryptDecryptRoundTripWithAES128Key(t *testing.T) { require.NoError(t, err) require.Equal(t, input, decrypted) } + +func TestDecryptDataRejectsValueWithoutEncryptionLayerHeader(t *testing.T) { + manager := NewEncryptionManager(&mockMetaManager{}) + + _, err := manager.DecryptData(context.Background(), 1, []byte("abc")) + require.Error(t, err) + require.Contains(t, err.Error(), "decryption failed") +} + +func TestDecryptDataRejectsInvalidUnencryptedHeader(t *testing.T) { + manager := NewEncryptionManager(&mockMetaManager{}) + + _, err := manager.DecryptData(context.Background(), 1, []byte{0x00, 0x01, 0x02, 0x03, 'x'}) + require.Error(t, err) + require.Contains(t, err.Error(), "decryption failed") +} + +func TestDecryptDataRejectsInvalidEncryptedHeader(t *testing.T) { + manager := NewEncryptionManager(&mockMetaManager{}) + + _, err := manager.DecryptData(context.Background(), 1, []byte{0x01, 0x00, 0x00, 0x00, 'x'}) + require.Error(t, err) + require.Contains(t, err.Error(), "decryption failed") +} diff --git a/pkg/encryption/format.go b/pkg/encryption/format.go index d10cbd366c..453b28bbdf 100644 --- a/pkg/encryption/format.go +++ b/pkg/encryption/format.go @@ -14,7 +14,9 @@ package encryption import ( + "github.com/pingcap/log" cerrors "github.com/pingcap/ticdc/pkg/errors" + "go.uber.org/zap" ) const ( @@ -125,25 +127,15 @@ func EncodeUnencryptedData(data []byte) []byte { return result } -// DecodeUnencryptedData decodes unencrypted data (removes header if present) +// DecodeUnencryptedData decodes unencrypted data by removing the 4-byte plaintext header. +// Callers must guarantee the value is marked as passing through the encryption layer +// and uses the unencrypted header format. func DecodeUnencryptedData(data []byte) ([]byte, error) { - if len(data) < EncryptionHeaderSize { - // No header, return as-is (backward compatibility) - return data, nil + if !HasUnencryptedHeader(data) { + log.Panic("unexpected data without unencrypted header", + zap.Int("dataLen", len(data))) } - - version := data[0] - if version == VersionUnencrypted && dataKeyIDIsZero(data) { - // New-format unencrypted data with header, remove header - return data[4:], nil - } - - // For backward compatibility, treat any other format as legacy unencrypted data - // This includes: - // - Legacy data without header (any pattern) - // - Data that might look like encrypted but is actually legacy - // The caller is responsible for ensuring data is not actually encrypted - return data, nil + return data[EncryptionHeaderSize:], nil } // ExtractDataKeyID extracts the data key ID from encrypted data diff --git a/pkg/encryption/format_test.go b/pkg/encryption/format_test.go index b89f4a2039..f4c6b951bd 100644 --- a/pkg/encryption/format_test.go +++ b/pkg/encryption/format_test.go @@ -139,31 +139,23 @@ func TestGetVersion(t *testing.T) { require.Equal(t, byte(0x00), GetVersion(shortData)) } -func TestDecodeUnencryptedDataBackwardCompatibility(t *testing.T) { - // Legacy data without header should be returned as-is - // Use data that is too short to have a header (length < 4) +func TestDecodeUnencryptedDataPanicsWithoutUnencryptedHeader(t *testing.T) { legacyData := []byte("legacy") - decoded, err := DecodeUnencryptedData(legacyData) - require.NoError(t, err) - require.Equal(t, legacyData, decoded) + require.Panics(t, func() { + _, _ = DecodeUnencryptedData(legacyData) + }) - // Also test with data that has non-zero DataKeyID pattern - // This can't be confused with new-format encrypted data (which would have non-zero key ID) - // and can't be confused with new-format unencrypted (which has zero key ID) legacyData2 := []byte{0x00, 0x01, 0x02, 0x03, 0x04, 0x05} - decoded2, err := DecodeUnencryptedData(legacyData2) - require.NoError(t, err) - require.Equal(t, legacyData2, decoded2) + require.Panics(t, func() { + _, _ = DecodeUnencryptedData(legacyData2) + }) } -func TestDecodeUnencryptedDataWithEncryptedData(t *testing.T) { - // For backward compatibility, DecodeUnencryptedData treats any format as legacy unencrypted data - // and returns the data as-is. It does not return an error even for encrypted-looking data. - // The caller is responsible for ensuring data is not actually encrypted. +func TestDecodeUnencryptedDataPanicsWithEncryptedData(t *testing.T) { encryptedData := []byte{0x01, 'a', 'b', 'c', 'd', 'a', 't', 'a'} - decoded, err := DecodeUnencryptedData(encryptedData) - require.NoError(t, err) - require.Equal(t, encryptedData, decoded) + require.Panics(t, func() { + _, _ = DecodeUnencryptedData(encryptedData) + }) } func TestExtractDataKeyID(t *testing.T) { From 3fa8a5b58b41633450221820b3501d784fa24b1b Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Thu, 7 May 2026 18:28:56 +0800 Subject: [PATCH 17/20] eventstore: clarify key boundary lengths Signed-off-by: tenfyzhong --- logservice/eventstore/format.go | 6 +++--- logservice/eventstore/format_test.go | 2 +- logservice/eventstore/pebble_test.go | 4 ++-- 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/logservice/eventstore/format.go b/logservice/eventstore/format.go index 7cec5b656d..29942031ef 100644 --- a/logservice/eventstore/format.go +++ b/logservice/eventstore/format.go @@ -81,13 +81,13 @@ const ( // encodeTxnCommitTsBoundaryKey encodes the event-store key boundary up to txnCommitTs. func encodeTxnCommitTsBoundaryKey(uniqueID uint64, tableID int64, txnCommitTs uint64) []byte { - buf := make([]byte, encodedKeyTxnStartTsOffset) + buf := make([]byte, encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen) encodeTxnCommitTsBoundaryKeyTo(buf, uniqueID, tableID, txnCommitTs) return buf } func encodeScanLowerBound(uniqueID uint64, tableID int64, txnCommitTs uint64, txnStartTs uint64) []byte { - buf := make([]byte, encodedKeyAttributesOffset) + buf := make([]byte, encodedKeyTxnStartTsOffset+encodedKeyTxnStartTsLen) encodeTxnCommitTsBoundaryKeyTo(buf, uniqueID, tableID, txnCommitTs) binary.BigEndian.PutUint64(buf[encodedKeyTxnStartTsOffset:encodedKeyTxnStartTsOffset+encodedKeyTxnStartTsLen], txnStartTs) return buf @@ -180,7 +180,7 @@ func KeyUsesEncryptionLayer(key []byte) bool { // It works for both full event keys and DeleteRange boundary keys because both // contain uniqueID, tableID, and txnCommitTs as the first three fields. func decodeTxnCommitTsFromEncodedKey(key []byte) (uint64, bool) { - if len(key) < encodedKeyTxnStartTsOffset { + if len(key) < encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen { return 0, false } return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsOffset : encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen]), true diff --git a/logservice/eventstore/format_test.go b/logservice/eventstore/format_test.go index 38bffa090e..71fcd14190 100644 --- a/logservice/eventstore/format_test.go +++ b/logservice/eventstore/format_test.go @@ -69,7 +69,7 @@ func TestEventStoreKeyFormatGolden(t *testing.T) { require.Equal(t, expectedKey[:encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen], encodeTxnCommitTsBoundaryKey(uniqueID, tableID, txnCommitTs)) - require.Equal(t, expectedKey[:encodedKeyAttributesOffset], + require.Equal(t, expectedKey[:encodedKeyTxnStartTsOffset+encodedKeyTxnStartTsLen], encodeScanLowerBound(uniqueID, tableID, txnCommitTs, txnStartTs)) decodedTxnCommitTs, ok := decodeTxnCommitTsFromEncodedKey(key) diff --git a/logservice/eventstore/pebble_test.go b/logservice/eventstore/pebble_test.go index 78a289d4aa..d143f9904b 100644 --- a/logservice/eventstore/pebble_test.go +++ b/logservice/eventstore/pebble_test.go @@ -163,11 +163,11 @@ func TestEventStoreKeyBounds(t *testing.T) { } key := EncodeKey(1, 1, event, CompressionNone) commitTsBoundaryKey := encodeTxnCommitTsBoundaryKey(1, 1, event.CRTs) - require.Len(t, commitTsBoundaryKey, encodedKeyTxnStartTsOffset) + require.Len(t, commitTsBoundaryKey, encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen) require.True(t, bytes.HasPrefix(key, commitTsBoundaryKey)) lowerBound := encodeScanLowerBound(1, 1, event.CRTs, event.StartTs) - require.Len(t, lowerBound, encodedKeyAttributesOffset) + require.Len(t, lowerBound, encodedKeyTxnStartTsOffset+encodedKeyTxnStartTsLen) require.True(t, bytes.HasPrefix(key, lowerBound)) previousEvent := &common.RawKVEntry{ From ab781e5bbcbf3b3fc7437675f6037805255fbc32 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Fri, 8 May 2026 10:52:12 +0800 Subject: [PATCH 18/20] logservice: remove unused eventstore raw value return Signed-off-by: tenfyzhong --- logservice/eventstore/event_store.go | 12 +++--- logservice/eventstore/event_store_test.go | 48 +++++++++++++++++++++++ 2 files changed, 54 insertions(+), 6 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 8f7d365349..38c81d30be 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1384,7 +1384,7 @@ func (e *eventStore) writeEvents( } } else if e.encryptionManager != nil { var value []byte - _, value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, needCompress) + value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, needCompress) valueBytesAfter = int64(len(value)) // Encrypt if encryption is enabled (after compression) @@ -1408,7 +1408,7 @@ func (e *eventStore) writeEvents( } } else { var value []byte - _, value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, true) + value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, true) valueBytesAfter = int64(len(value)) // SetDeferred is a write path optimization. Now that the compressed // value length is known, reserve the exact key/value space in the @@ -1458,15 +1458,15 @@ func encodeAndMaybeCompressValue( rawBuf []byte, dstBuf []byte, needCompress bool, -) (rawValue []byte, value []byte, compressionType CompressionType, nextRawBuf []byte, nextDstBuf []byte) { - rawValue = ensureValueBuffer(rawBuf, int(kv.GetSize())) +) (value []byte, compressionType CompressionType, nextRawBuf []byte, nextDstBuf []byte) { + rawValue := ensureValueBuffer(rawBuf, int(kv.GetSize())) rawValue = kv.EncodeTo(rawValue) value = rawValue compressionType = CompressionNone nextRawBuf = rawValue[:0] nextDstBuf = dstBuf if !needCompress { - return rawValue, value, compressionType, nextRawBuf, nextDstBuf + return value, compressionType, nextRawBuf, nextDstBuf } maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) @@ -1475,7 +1475,7 @@ func encodeAndMaybeCompressValue( compressionType = CompressionZSTD nextDstBuf = value[:0] metrics.EventStoreCompressedRowsCount.Inc() - return rawValue, value, compressionType, nextRawBuf, nextDstBuf + return value, compressionType, nextRawBuf, nextDstBuf } func ensureValueBuffer(buf []byte, minCap int) []byte { diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index a5b81baeca..db28856511 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -1203,6 +1203,54 @@ func TestWriteToEventStoreZstdCompressionDisabled(t *testing.T) { require.Equal(t, 1, count) } +func TestEncodeAndMaybeCompressValue(t *testing.T) { + entry := &common.RawKVEntry{ + OpType: common.OpTypePut, + StartTs: 100, + CRTs: 200, + Key: []byte("encode-key"), + Value: bytes.Repeat([]byte("value"), 32), + } + + encoder, err := zstd.NewWriter(nil) + require.NoError(t, err) + defer encoder.Close() + + expectedRawValue := entry.Encode() + + t.Run("withoutCompression", func(t *testing.T) { + dstBuf := make([]byte, 0, 8) + value, compressionType, nextRawBuf, nextDstBuf := encodeAndMaybeCompressValue( + entry, encoder, nil, dstBuf, false, + ) + require.Equal(t, CompressionNone, compressionType) + require.Equal(t, expectedRawValue, value) + require.Len(t, nextRawBuf, 0) + require.GreaterOrEqual(t, cap(nextRawBuf), len(expectedRawValue)) + require.Empty(t, nextDstBuf) + require.Equal(t, cap(dstBuf), cap(nextDstBuf)) + }) + + t.Run("withCompression", func(t *testing.T) { + value, compressionType, nextRawBuf, nextDstBuf := encodeAndMaybeCompressValue( + entry, encoder, nil, nil, true, + ) + require.Equal(t, CompressionZSTD, compressionType) + require.Len(t, nextRawBuf, 0) + require.GreaterOrEqual(t, cap(nextRawBuf), len(expectedRawValue)) + require.Empty(t, nextDstBuf) + require.GreaterOrEqual(t, cap(nextDstBuf), len(value)) + + decoder, err := zstd.NewReader(nil) + require.NoError(t, err) + defer decoder.Close() + + decodedValue, err := decoder.DecodeAll(value, nil) + require.NoError(t, err) + require.Equal(t, expectedRawValue, decodedValue) + }) +} + func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { restoreCfg := setZstdCompressionForTest(t, true) defer restoreCfg() From a43786ad78c82f9065ee4c25334f8b4c77e5e198 Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Sat, 9 May 2026 10:59:44 +0800 Subject: [PATCH 19/20] pkg/encryption: remove decode panic tests Signed-off-by: tenfyzhong --- pkg/encryption/format_test.go | 19 ------------------- 1 file changed, 19 deletions(-) diff --git a/pkg/encryption/format_test.go b/pkg/encryption/format_test.go index f4c6b951bd..e10d1553fc 100644 --- a/pkg/encryption/format_test.go +++ b/pkg/encryption/format_test.go @@ -139,25 +139,6 @@ func TestGetVersion(t *testing.T) { require.Equal(t, byte(0x00), GetVersion(shortData)) } -func TestDecodeUnencryptedDataPanicsWithoutUnencryptedHeader(t *testing.T) { - legacyData := []byte("legacy") - require.Panics(t, func() { - _, _ = DecodeUnencryptedData(legacyData) - }) - - legacyData2 := []byte{0x00, 0x01, 0x02, 0x03, 0x04, 0x05} - require.Panics(t, func() { - _, _ = DecodeUnencryptedData(legacyData2) - }) -} - -func TestDecodeUnencryptedDataPanicsWithEncryptedData(t *testing.T) { - encryptedData := []byte{0x01, 'a', 'b', 'c', 'd', 'a', 't', 'a'} - require.Panics(t, func() { - _, _ = DecodeUnencryptedData(encryptedData) - }) -} - func TestExtractDataKeyID(t *testing.T) { data := []byte("payload") keyID := "xyz" From 78b971f8aa7fd3333b8d783e53c04e161a593c1c Mon Sep 17 00:00:00 2001 From: tenfyzhong Date: Mon, 11 May 2026 17:24:06 +0800 Subject: [PATCH 20/20] tests: move debezium_basic to heavy integration tests Signed-off-by: tenfyzhong --- tests/integration_tests/run_heavy_it_in_ci.sh | 4 ++-- tests/integration_tests/run_light_it_in_ci.sh | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/tests/integration_tests/run_heavy_it_in_ci.sh b/tests/integration_tests/run_heavy_it_in_ci.sh index 7ae9105f96..4209ce706f 100755 --- a/tests/integration_tests/run_heavy_it_in_ci.sh +++ b/tests/integration_tests/run_heavy_it_in_ci.sh @@ -95,7 +95,7 @@ kafka_groups=( # G13 'debezium01 fail_over_ddl_mix' # G14 - 'debezium02' + 'debezium_basic debezium02' # G15 'debezium03' ) @@ -131,7 +131,7 @@ pulsar_groups=( # G13 'debezium01 fail_over_ddl_mix' # G14 - 'debezium02' + 'debezium_basic debezium02' # G15 'debezium03' ) diff --git a/tests/integration_tests/run_light_it_in_ci.sh b/tests/integration_tests/run_light_it_in_ci.sh index b70548517a..a2fea1ccc3 100755 --- a/tests/integration_tests/run_light_it_in_ci.sh +++ b/tests/integration_tests/run_light_it_in_ci.sh @@ -99,7 +99,7 @@ kafka_groups=( # G13 'cli_tls_with_auth cli_with_auth fail_over_ddl_N maintainer_failover_when_operator' # G14 - 'kafka_simple_basic avro_basic debezium_basic fail_over_ddl_O update_changefeed_check_config' + 'kafka_simple_basic avro_basic fail_over_ddl_O update_changefeed_check_config' # G15 'kafka_simple_basic_avro split_region autorandom gc_safepoint kafka_log_info' ) @@ -137,7 +137,7 @@ pulsar_groups=( # G13 'cli_tls_with_auth cli_with_auth fail_over_ddl_N maintainer_failover_when_operator' # G14 - 'avro_basic debezium_basic fail_over_ddl_O update_changefeed_check_config' + 'avro_basic fail_over_ddl_O update_changefeed_check_config' # G15 'split_region autorandom gc_safepoint' )