feat(datanode): Enable crc check for BaseExtentID information #23118077

Signed-off-by: leonrayang <chl696@sina.com>
This commit is contained in:
leonrayang 2025-03-06 15:53:37 +08:00 committed by Victor1319
parent b1672c4e5d
commit 69e7bc5475
5 changed files with 139 additions and 25 deletions

View File

@ -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)
}

View File

@ -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
}

View File

@ -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) {

View File

@ -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)

View File

@ -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 {