diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 4138cb3609..1a7ef473bf 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -35,6 +35,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" @@ -118,6 +119,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, @@ -184,9 +187,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 @@ -242,6 +246,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 ( @@ -262,6 +268,8 @@ 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) dbs, pebbleCache, tableCache := createPebbleDBs(dbPath, dbCount) store := &eventStore{ pdClock: appcontext.GetService[pdutil.Clock](appcontext.DefaultPDClock), @@ -287,6 +295,7 @@ func New( }, compressionThreshold: config.GetGlobalServerConfig().Debug.EventStore.CompressionThreshold, enableZstdCompression: config.GetGlobalServerConfig().Debug.EventStore.EnableZstdCompression, + encryptionManager: encMgr, } store.gcManager = newGCManager(store.dbs, deleteDataRange, compactDataRange) @@ -501,6 +510,7 @@ func (e *eventStore) RegisterDispatcher( dispatcherID: dispatcherID, tableSpan: dispatcherSpan, checkpointTs: startTs, + keyspaceID: dispatcherSpan.KeyspaceID, } stat.resolvedTs.Store(startTs) @@ -627,6 +637,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(), @@ -951,16 +962,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, }, nil } @@ -1364,37 +1377,62 @@ func (e *eventStore) writeEvents( continue } - compressionType := CompressionNone 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)) + 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) { + 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 { + var value []byte + value, compressionType, rawBuf, dstBuf = encodeAndMaybeCompressValue(kv, encoder, rawBuf, dstBuf, needCompress) + valueBytesAfter = int64(len(value)) - if e.enableZstdCompression && valueBytesBefore > int64(e.compressionThreshold) { - if cap(rawBuf) < int(valueBytesBefore) { - rawBuf = make([]byte, 0, int(valueBytesBefore)) - } else { - rawBuf = rawBuf[:0] + // Encrypt if encryption is enabled (after compression) + encryptedValue, err := e.encryptionManager.EncryptData(context.Background(), event.keyspaceID, value) + if err != nil { + return err } - rawValue := kv.EncodeTo(rawBuf) - maxEncodedSize := encoder.MaxEncodedSize(len(rawValue)) - if cap(dstBuf) < maxEncodedSize { - dstBuf = make([]byte, 0, maxEncodedSize) - } else { - dstBuf = dstBuf[:0] + + op := batch.SetDeferred(keyLen, len(encryptedValue)) + if err := encodeDeferredEventKey(op, keyLen, uint64(event.subID), event.tableID, kv, compressionType, true); err != nil { + return err } - value := encoder.EncodeAll(rawValue, dstBuf) + 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 + } + } else { + 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] @@ -1405,28 +1443,7 @@ func (e *eventStore) writeEvents( if err := op.Finish(); err != nil { return err } - 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 } @@ -1451,6 +1468,60 @@ func (e *eventStore) writeEvents( return err } +func encodeAndMaybeCompressValue( + kv *common.RawKVEntry, + encoder *zstd.Encoder, + rawBuf []byte, + dstBuf []byte, + needCompress bool, +) (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 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 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` @@ -1468,6 +1539,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) { @@ -1479,6 +1553,17 @@ func (iter *eventStoreIter) Next() (*common.RawKVEntry, bool) { key := iter.innerIter.Key() value := iter.innerIter.Value() + 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)) + } + value = decryptedValue + } _, compressionType := DecodeKeyAttributes(key) var decodedValue []byte if compressionType == CompressionZSTD { diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index 5b3a8ac9ec..61ba0ee9e9 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,8 @@ 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" + 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" @@ -48,6 +51,50 @@ type mockSubscriptionClient struct { subscriptions map[logpuller.SubscriptionID]*mockSubscriptionStat } +type spyEncryptionManager struct { + encryptKeyspaceID uint32 + decryptKeyspaceID uint32 + encryptCalls int + 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++ + encrypted := make([]byte, encryption.EncryptionHeaderSize+len(data)) + encrypted[0] = 0x01 + encrypted[3] = 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), @@ -190,6 +237,187 @@ 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, err := es.GetIterator(dispatcherID, dataRange) + require.NoError(t, err) + 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 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, err := es.GetIterator(dispatcherID, common.DataRange{ + Span: span, + CommitTsStart: 0, + CommitTsEnd: largeKV.CRTs, + }) + require.NoError(t, err) + 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] @@ -501,7 +729,7 @@ func TestGetIteratorPanicWhenStartLessThanCheckpoint(t *testing.T) { store.UpdateDispatcherCheckpointTs(dispatcherID, 120) require.Panics(t, func() { - store.GetIterator(dispatcherID, common.DataRange{ + _, _ = store.GetIterator(dispatcherID, common.DataRange{ Span: span, CommitTsStart: 110, CommitTsEnd: 150, @@ -998,6 +1226,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() @@ -1086,6 +1362,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 47cd09eb74..29942031ef 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,53 +39,82 @@ const ( ) const ( - encodedKeyUint64Len = 8 - encodedKeyTxnCommitTsStart = 2 * encodedKeyUint64Len - encodedKeyTxnCommitTsEnd = encodedKeyTxnCommitTsStart + encodedKeyUint64Len - encodedKeyAttributesOffset = 4 * encodedKeyUint64Len - encodedKeyAttributesEnd = encodedKeyAttributesOffset + 2 + // 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 + 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 = encodedKeyAttributesOffset + encodedKeyAttributesLen ) 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 + // encodedKeyMaskEncryption indicates the value passed through the encryption layer. + encodedKeyMaskEncryption uint64 = 1 << iota ) // 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, 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[encodedKeyTxnCommitTsEnd:encodedKeyAttributesOffset], 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[:encodedKeyUint64Len], uniqueID) - binary.BigEndian.PutUint64(buf[encodedKeyUint64Len:encodedKeyTxnCommitTsStart], uint64(tableID)) - binary.BigEndian.PutUint64(buf[encodedKeyTxnCommitTsStart:encodedKeyTxnCommitTsEnd], 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 { - // uniqueID, tableID, txnCommitTs, txnStartTs, Put/Delete, CompressionType, Key + // uniqueID, tableID, txnCommitTs, txnStartTs, DMLOrder, CompressionType, Mask, Key return encodedKeyAttributesEnd + len(event.Key) } -// EncodeKeyTo appends an encoded event-store key to buf. -// Format: uniqueID, tableID, txnCommitTs, txnStartTs, delete/update/insert, Key. -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, tableID int64, event *common.RawKVEntry, compressionType CompressionType, + usesEncryptionLayer bool, ) []byte { if event == nil { log.Panic("rawkv must not be nil", zap.Any("event", event)) @@ -100,12 +129,38 @@ func EncodeKeyTo( buf = binary.BigEndian.AppendUint64(buf, event.StartTs) // Let Delete < Update < Insert dmlOrder := getDMLOrder(event) - combinedOrder := uint16(compressionType) | (uint16(dmlOrder) << dmlOrderShift) - buf = binary.BigEndian.AppendUint16(buf, combinedOrder) + buf = append(buf, byte(dmlOrder)) + buf = append(buf, byte(compressionType)) + mask := uint64(0) + if usesEncryptionLayer { + mask |= encodedKeyMaskEncryption + } + buf = binary.BigEndian.AppendUint64(buf, mask) // key return append(buf, event.Key...) } +// EncodeKeyTo appends an encoded event-store key to buf. +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) @@ -113,18 +168,22 @@ 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]) - return DMLOrder((combinedOrder & dmlOrderMask) >> dmlOrderShift), CompressionType(combinedOrder & compressionMask) + return DMLOrder(key[encodedKeyDMLOrderOffset]), CompressionType(key[encodedKeyCompressionOffset]) +} + +func KeyUsesEncryptionLayer(key []byte) bool { + mask := binary.BigEndian.Uint64(key[encodedKeyMaskOffset : encodedKeyMaskOffset+encodedKeyMaskLen]) + return mask&encodedKeyMaskEncryption != 0 } // decodeTxnCommitTsFromEncodedKey decodes txnCommitTs from an event-store key boundary. // 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) < encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen { return 0, false } - return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsStart:encodedKeyTxnCommitTsEnd]), true + return binary.BigEndian.Uint64(key[encodedKeyTxnCommitTsOffset : encodedKeyTxnCommitTsOffset+encodedKeyTxnCommitTsLen]), true } // getDMLOrder returns the order of the dml types: delete