From 69e7bc54757a3776f924986267dc2d71b870fc11 Mon Sep 17 00:00:00 2001 From: leonrayang Date: Thu, 6 Mar 2025 15:53:37 +0800 Subject: [PATCH] feat(datanode): Enable crc check for BaseExtentID information #23118077 Signed-off-by: leonrayang --- datanode/partition.go | 9 ++ datanode/storage/extent_store.go | 22 ++++- datanode/storage/extent_store_test.go | 16 ++++ datanode/storage/persistence_crc.go | 115 +++++++++++++++++++++----- master/master_quota_manager.go | 2 +- 5 files changed, 139 insertions(+), 25 deletions(-) diff --git a/datanode/partition.go b/datanode/partition.go index 7c7d3d594..38bf2eb88 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -451,6 +451,7 @@ func newDataPartition(dpCfg *dataPartitionCfg, disk *Disk, isCreate bool) (dp *D partition.replicasInit() partition.extentStore, err = storage.NewExtentStore(partition.path, dpCfg.PartitionID, dpCfg.PartitionSize, partition.partitionType, disk.dataNode.cacheCap, isCreate) + partition.extentStore.IsEnableSnapshot = dpCfg.IsEnableSnapshot if err != nil { log.LogWarnf("action[newDataPartition] dp %v NewExtentStore failed %v", partitionID, err.Error()) return @@ -897,6 +898,9 @@ func (dp *DataPartition) statusUpdate() { if dp.isNormalType() && dp.extentStore.GetExtentCount() >= storage.MaxExtentCount { status = proto.ReadOnly } + if dp.IsBaseFileIDException() { + status = proto.ReadOnly + } if dp.disk.Status == proto.ReadOnly { status = proto.ReadOnly } @@ -1337,6 +1341,11 @@ func (dp *DataPartition) putRepairConn(conn net.Conn, forceClose bool) { dp.dataNode.putRepairConnFunc(conn, forceClose) } +func (dp *DataPartition) IsBaseFileIDException() bool { + extID, _ := dp.extentStore.GetPersistenceBaseExtentID() + return extID > storage.MaxExtentID +} + func (dp *DataPartition) isNormalType() bool { return proto.IsNormalDp(dp.partitionType) } diff --git a/datanode/storage/extent_store.go b/datanode/storage/extent_store.go index 5938ace92..6a64c0adf 100644 --- a/datanode/storage/extent_store.go +++ b/datanode/storage/extent_store.go @@ -58,6 +58,7 @@ const ( RandomWriteType = 2 AppendWriteType = 1 AppendRandomWriteType = 4 + MaxExtentID = 1 << 30 NormalExtentDeleteRetainTime = 3600 * 4 CacheFlushMinInterval = 5 * time.Minute @@ -151,6 +152,8 @@ type ExtentStore struct { stopC chan interface{} ApplyId uint64 DirectRead bool + IsEnableSnapshot bool + extIDLock sync.Mutex } func MkdirAll(name string) (err error) { @@ -575,11 +578,9 @@ func (s *ExtentStore) initBaseFileID() error { defer func() { log.LogInfof("[initBaseFileID] store(%v) init base file id using time(%v), count(%v)", s.dataPath, time.Since(begin), extNum) }() - var baseFileID uint64 - baseFileID, _ = s.GetPersistenceBaseExtentID() - log.LogInfof("[initBaseFileID] store(%v) init base file to persistence base extent id using time(%v)", s.dataPath, time.Since(begin)) // NOTE: try to read hint + var baseFileID uint64 var err error var extMap map[uint64]*ExtentInfo extMap, err = s.readReadDirHint() @@ -633,6 +634,20 @@ func (s *ExtentStore) initBaseFileID() error { if baseFileID < MinExtentID { baseFileID = MinExtentID } + + diskBaseFileID, _ := s.GetPersistenceBaseExtentID() + log.LogInfof("[initBaseFileID] store(%v) init base file to persistence base extent id using time(%v)", s.dataPath, time.Since(begin)) + + if errTmp := s.CheckBaseExtentCrc(); errTmp != nil && errTmp != io.EOF { + log.LogErrorf("[initBaseFileID] store(%v) init base file but not consistent baseFileID %v diskBaseFileID %v err %v", s.dataPath, baseFileID, diskBaseFileID, errTmp) + return errTmp + } + if baseFileID != diskBaseFileID { + log.LogWarnf("[initBaseFileID] store(%v) init base file but not consistent baseFileID %v diskBaseFileID %v", s.dataPath, baseFileID, diskBaseFileID) + if diskBaseFileID > baseFileID { + baseFileID = diskBaseFileID + } + } atomic.StoreUint64(&s.baseExtentID, baseFileID) log.LogInfof("datadir(%v) maxBaseId(%v)", s.dataPath, baseFileID) return nil @@ -1399,7 +1414,6 @@ func (s *ExtentStore) UpdateBaseExtentID(id uint64) (err error) { err = s.PersistenceBaseExtentID(atomic.LoadUint64(&s.baseExtentID)) } s.PreAllocSpaceOnVerfiyFile(atomic.LoadUint64(&s.baseExtentID)) - return } diff --git a/datanode/storage/extent_store_test.go b/datanode/storage/extent_store_test.go index 0dc2eced0..9c68be1d5 100644 --- a/datanode/storage/extent_store_test.go +++ b/datanode/storage/extent_store_test.go @@ -272,6 +272,19 @@ func staleExtentStoreTest(t *testing.T, dpType int) { defer newS2.Close() } +func extentStoreBaseExtentTest(t *testing.T, s *storage.ExtentStore) { + normalId, err := s.NextExtentID() + require.NoError(t, err) + s.WritePreAllocSpaceExtentIDOnVerifyFile(normalId + 1000) + err = s.CheckBaseExtentCrc() + require.NoError(t, err) + + fNormalId, _ := s.GetPersistenceBaseExtentID() + fSpaceNormalId := s.GetPreAllocSpaceExtentIDOnVerifyFile() + require.Equal(t, fNormalId, normalId) + require.Equal(t, fSpaceNormalId, normalId+1000) +} + func ExtentStoreTest(t *testing.T, dpType int) { path, clean, err := getTestPathExtentStore() require.NoError(t, err) @@ -282,6 +295,9 @@ func ExtentStoreTest(t *testing.T, dpType int) { extentStoreLogicalTest(t, s) reopenExtentStoreTest(t, dpType) staleExtentStoreTest(t, dpType) + if dpType == proto.PartitionTypeNormal { + extentStoreBaseExtentTest(t, s) + } } func TestExtentStores(t *testing.T) { diff --git a/datanode/storage/persistence_crc.go b/datanode/storage/persistence_crc.go index 6e06777ee..1661fa925 100644 --- a/datanode/storage/persistence_crc.go +++ b/datanode/storage/persistence_crc.go @@ -16,6 +16,8 @@ package storage import ( "encoding/binary" + "fmt" + "hash/crc32" "io" "os" "path" @@ -34,7 +36,9 @@ type BlockCrc struct { type BlockCrcArr []*BlockCrc const ( - BaseExtentIDOffset = 0 + BaseExtentIDOffset = 0 + BaseExtentEndIDOffset = 8 + BaseExtentCrcOffset = 16 ) func (arr BlockCrcArr) Len() int { return len(arr) } @@ -137,16 +141,98 @@ func (s *ExtentStore) DeleteBlockCrc(extentID uint64) (err error) { return } -func (s *ExtentStore) PersistenceBaseExtentID(extentID uint64) (err error) { - value := make([]byte, 8) - binary.BigEndian.PutUint64(value, extentID) - _, err = s.metadataFp.WriteAt(value, BaseExtentIDOffset) +func (s *ExtentStore) calcExtentCrc() (crc uint32, err error) { + data := make([]byte, 16) + _, err = s.metadataFp.ReadAt(data, 0) + if err != nil { + return + } + + sign := crc32.NewIEEE() + if _, err = sign.Write(data); err != nil { + return + } + crc = sign.Sum32() return } -func (s *ExtentStore) GetPreAllocSpaceExtentIDOnVerifyFile() (extentID uint64) { +func (s *ExtentStore) PersistentBaseExtentCrc() (err error) { + var crc uint32 + if crc, err = s.calcExtentCrc(); err != nil { + if err != io.EOF { + return + } + return nil + } + dataCrc := make([]byte, 8) + binary.BigEndian.PutUint32(dataCrc, crc) + _, err = s.metadataFp.WriteAt(dataCrc, BaseExtentCrcOffset) + return +} + +func (s *ExtentStore) CheckBaseExtentCrc() (err error) { + var ( + crcCalc uint32 + crcRead uint32 + ) + if crcCalc, err = s.calcExtentCrc(); err != nil { + if err != io.EOF { + log.LogErrorf("CheckBaseExtentCrc dp %v err %v", s.partitionID, err) + } + return + } + data := make([]byte, 4) + if _, err = s.metadataFp.ReadAt(data, BaseExtentCrcOffset); err == io.EOF { + return nil // not init before + } + crcRead = binary.BigEndian.Uint32(data) + if crcRead != crcCalc { + err = fmt.Errorf("CheckBaseExtentCrc dp %v crc not equal %v vs %v", s.partitionID, crcRead, crcCalc) + } + return +} + +func (s *ExtentStore) PersistenceBaseExtentID(extentID uint64) (err error) { + s.extIDLock.Lock() + defer s.extIDLock.Unlock() + value := make([]byte, 8) - _, err := s.metadataFp.ReadAt(value, 8) + binary.BigEndian.PutUint64(value, extentID) + _, err = s.metadataFp.WriteAt(value, BaseExtentIDOffset) + + return s.PersistentBaseExtentCrc() +} + +func (s *ExtentStore) GetPersistenceBaseExtentID() (extentID uint64, err error) { + s.extIDLock.Lock() + defer s.extIDLock.Unlock() + + data := make([]byte, 8) + _, err = s.metadataFp.ReadAt(data, 0) + if err != nil { + return + } + extentID = binary.BigEndian.Uint64(data) + return +} + +func (s *ExtentStore) WritePreAllocSpaceExtentIDOnVerifyFile(extentID uint64) (err error) { + s.extIDLock.Lock() + defer s.extIDLock.Unlock() + + value := make([]byte, 8) + binary.BigEndian.PutUint64(value, extentID) + _, err = s.metadataFp.WriteAt(value, BaseExtentEndIDOffset) + + return s.PersistentBaseExtentCrc() +} + +func (s *ExtentStore) GetPreAllocSpaceExtentIDOnVerifyFile() (extentID uint64) { + s.extIDLock.Lock() + defer s.extIDLock.Unlock() + + value := make([]byte, 8) + _, err := s.metadataFp.ReadAt(value, BaseExtentEndIDOffset) if err != nil { return } @@ -197,11 +283,10 @@ func (s *ExtentStore) PreAllocSpaceOnVerfiyFile(currExtentID uint64) { } } - data := make([]byte, 8) - binary.BigEndian.PutUint64(data, uint64(endAllocSpaceExtentID)) - if _, err = s.metadataFp.WriteAt(data, 8); err != nil { + if err = s.WritePreAllocSpaceExtentIDOnVerifyFile(uint64(endAllocSpaceExtentID)); err != nil { return } + atomic.StoreUint64(&s.hasAllocSpaceExtentIDOnVerfiyFile, uint64(endAllocSpaceExtentID)) log.LogInfof("Action(PreAllocSpaceOnVerifyFile) PartitionID(%v) currentExtent(%v)"+ "PrevAllocSpaceExtentIDOnVerifyFile(%v) EndAllocSpaceExtentIDOnVerifyFile(%v)"+ @@ -210,16 +295,6 @@ func (s *ExtentStore) PreAllocSpaceOnVerfiyFile(currExtentID uint64) { } } -func (s *ExtentStore) GetPersistenceBaseExtentID() (extentID uint64, err error) { - data := make([]byte, 8) - _, err = s.metadataFp.ReadAt(data, 0) - if err != nil { - return - } - extentID = binary.BigEndian.Uint64(data) - return -} - func (s *ExtentStore) PersistenceHasDeleteExtent(extentID uint64) (err error) { data := make([]byte, 8) binary.BigEndian.PutUint64(data, extentID) diff --git a/master/master_quota_manager.go b/master/master_quota_manager.go index 60aaa7638..e5a9fbf99 100644 --- a/master/master_quota_manager.go +++ b/master/master_quota_manager.go @@ -43,7 +43,7 @@ func (mqMgr *MasterQuotaManager) persistQuota(quotaInfo *proto.QuotaInfo) (err e metadata := new(RaftCmd) metadata.Op = opSyncSetQuota - metadata.K = quotaPrefix + strconv.FormatUint(mqMgr.vol.ID, 10) + keySeparator + strconv.FormatUint(uint64(quotaId), 10) + metadata.K = quotaPrefix + strconv.FormatUint(mqMgr.vol.ID, 10) + keySeparator + strconv.FormatUint(uint64(quotaInfo.QuotaId), 10) metadata.V = value if err = mqMgr.c.submit(metadata); err != nil {