diff --git a/datanode/data_partition_repair.go b/datanode/data_partition_repair.go index 2667db506..048aeae26 100644 --- a/datanode/data_partition_repair.go +++ b/datanode/data_partition_repair.go @@ -843,7 +843,19 @@ func (dp *DataPartition) streamRepairExtent(remoteExtentInfo *storage.ExtentInfo } } else { log.LogDebugf("streamRepairExtent reply size %v, currFixoffset %v, reply %v ", reply.GetSize(), currFixOffset, reply) - _, err = store.Write(uint64(localExtentInfo.FileID), int64(currFixOffset), int64(reply.GetSize()), reply.GetData(), reply.GetCRC(), wType, BufferWrite, isEmptyResponse, true, request.GetOpcode() == proto.OpBackupWrite) + param := &storage.WriteParam{ + ExtentID: uint64(localExtentInfo.FileID), + Offset: int64(currFixOffset), + Size: int64(reply.GetSize()), + Data: reply.GetData(), + Crc: reply.GetCRC(), + WriteType: wType, + IsSync: BufferWrite, + IsHole: isEmptyResponse, + IsRepair: true, + IsBackupWrite: request.GetOpcode() == proto.OpBackupWrite, + } + _, err = store.Write(param) } // log.LogDebugf("streamRepairExtent reply size %v, currFixoffset %v, reply %v err %v", reply.Size, currFixOffset, reply, err) // write to the local extent file diff --git a/datanode/data_partition_repair_test.go b/datanode/data_partition_repair_test.go index ef3f2b873..313600e54 100644 --- a/datanode/data_partition_repair_test.go +++ b/datanode/data_partition_repair_test.go @@ -262,13 +262,26 @@ func mockMakeDp(path string) *DataPartition { func extentStoreNormalRwTest(t *testing.T, s *storage.ExtentStore, id uint64, crc uint32, data []byte) { // append write - _, err := s.Write(id, 0, int64(len(data)), data, crc, storage.AppendWriteType, true, false, false, false) + param := &storage.WriteParam{ + ExtentID: id, + Offset: 0, + Size: int64(len(data)), + Data: data, + Crc: crc, + WriteType: storage.AppendWriteType, + IsSync: true, + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err := s.Write(param) require.NoError(t, err) actualCrc, err := s.Read(id, 0, int64(len(data)), data, false, false) require.NoError(t, err) require.EqualValues(t, crc, actualCrc) // random write - _, err = s.Write(id, 0, int64(len(data)), data, crc, storage.RandomWriteType, true, false, false, false) + param.WriteType = storage.RandomWriteType + _, err = s.Write(param) require.NoError(t, err) actualCrc, err = s.Read(id, 0, int64(len(data)), data, false, false) require.NoError(t, err) @@ -289,20 +302,39 @@ func extentReloadCheckNormalCrc(t *testing.T, s *storage.ExtentStore, id uint64, func extentStoreSnapshotRwTest(t *testing.T, s *storage.ExtentStore, id uint64, crc uint32, data []byte) { // append write + param := &storage.WriteParam{ + ExtentID: id, + Offset: int64(util.ExtentSize), + Size: int64(len(data)), + Data: data, + Crc: crc, + IsSync: true, + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + offset := int64(util.ExtentSize) - _, err := s.Write(id, offset, int64(len(data)), data, crc, storage.AppendRandomWriteType, true, false, false, false) + param.WriteType = storage.AppendRandomWriteType + _, err := s.Write(param) require.NoError(t, err) - _, err = s.Write(id, 0, int64(len(data)), data, crc, storage.AppendRandomWriteType, true, false, false, false) + param.WriteType = storage.AppendRandomWriteType + param.Offset = 0 + _, err = s.Write(param) assert.True(t, err != nil) - _, err = s.Write(id, offset, int64(len(data)), data, crc, storage.AppendWriteType, true, false, false, false) + + param.WriteType = storage.AppendWriteType + param.Offset = offset + _, err = s.Write(param) assert.True(t, err != nil) actualCrc, err := s.Read(id, offset, int64(len(data)), data, false, false) require.NoError(t, err) require.EqualValues(t, crc, actualCrc) // random write - _, err = s.Write(id, offset, int64(len(data)), data, crc, storage.RandomWriteType, true, false, false, false) + param.WriteType = storage.RandomWriteType + _, err = s.Write(param) require.NoError(t, err) actualCrc, err = s.Read(id, offset, int64(len(data)), data, false, false) require.NoError(t, err) @@ -310,7 +342,8 @@ func extentStoreSnapshotRwTest(t *testing.T, s *storage.ExtentStore, id uint64, // TODO: append random write require.NotEqualValues(t, s.GetStoreUsedSize(), 0) - _, err = s.Write(id, offset, int64(len(data)), data, crc, storage.AppendRandomWriteType, true, false, false, false) + param.WriteType = storage.AppendRandomWriteType + _, err = s.Write(param) require.NoError(t, err) // extent crc check @@ -321,15 +354,15 @@ func extentStoreSnapshotRwTest(t *testing.T, s *storage.ExtentStore, id uint64, assert.True(t, crc == extCrc) // check - offset = int64(util.ExtentSize)*2 + util.BlockSize - _, err = s.Write(id, offset, int64(len(data)), data, crc, storage.AppendRandomWriteType, true, false, false, false) + param.Offset = int64(util.ExtentSize)*2 + util.BlockSize + _, err = s.Write(param) require.NoError(t, err) e, err = s.LoadExtentFromDisk(id, true) assert.True(t, err == nil) - extCrc = e.GetCrc(offset / util.BlockSize) + extCrc = e.GetCrc(param.Offset / util.BlockSize) assert.True(t, crc == extCrc) - extCrc = e.GetCrc(offset/util.BlockSize + 1) + extCrc = e.GetCrc(param.Offset/util.BlockSize + 1) assert.True(t, 0 == extCrc) } diff --git a/datanode/partition_op_by_raft.go b/datanode/partition_op_by_raft.go index e419fd60a..a09a4d32c 100644 --- a/datanode/partition_op_by_raft.go +++ b/datanode/partition_op_by_raft.go @@ -282,7 +282,19 @@ func (dp *DataPartition) ApplyRandomWrite(command []byte, raftApplyID uint64) (r } dp.disk.limitWrite.Run(int(opItem.size), func() { - respStatus, err = dp.ExtentStore().Write(opItem.extentID, opItem.offset, opItem.size, opItem.data, opItem.crc, writeType, syncWrite, false, false, opItem.opcode == proto.OpBackupWrite) + param := &storage.WriteParam{ + ExtentID: uint64(opItem.extentID), + Offset: int64(opItem.offset), + Size: opItem.size, + Data: opItem.data, + Crc: opItem.crc, + WriteType: writeType, + IsSync: syncWrite, + IsHole: false, + IsRepair: false, + IsBackupWrite: opItem.opcode == proto.OpBackupWrite, + } + respStatus, err = dp.ExtentStore().Write(param) }) if err == nil || err == storage.ErrStoreAlreadyClosed { break diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index daa082f92..679bfe713 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -824,7 +824,19 @@ func (s *DataNode) handleWritePacket(p *repl.Packet) { partition.disk.allocCheckLimit(proto.IopsWriteType, 1) if writable := partition.disk.limitWrite.TryRun(int(p.Size), func() { - _, err = store.Write(p.ExtentID, p.ExtentOffset, int64(p.Size), p.Data, p.CRC, storage.AppendWriteType, p.IsSyncWrite(), false, false, false) + param := &storage.WriteParam{ + ExtentID: p.ExtentID, + Offset: p.ExtentOffset, + Size: int64(p.Size), + Data: p.Data, + Crc: p.CRC, + WriteType: storage.AppendWriteType, + IsSync: p.IsSyncWrite(), + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err = store.Write(param) }); !writable { err = storage.LimitedIoError return @@ -846,7 +858,19 @@ func (s *DataNode) handleWritePacket(p *repl.Packet) { partition.disk.allocCheckLimit(proto.IopsWriteType, 1) if writable := partition.disk.limitWrite.TryRun(int(p.Size), func() { - _, err = store.Write(p.ExtentID, p.ExtentOffset, int64(p.Size), p.Data, p.CRC, storage.AppendWriteType, p.IsSyncWrite(), false, false, false) + param := &storage.WriteParam{ + ExtentID: p.ExtentID, + Offset: p.ExtentOffset, + Size: int64(p.Size), + Data: p.Data, + Crc: p.CRC, + WriteType: storage.AppendWriteType, + IsSync: p.IsSyncWrite(), + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err = store.Write(param) }); !writable { err = storage.LimitedIoError return @@ -874,7 +898,19 @@ func (s *DataNode) handleWritePacket(p *repl.Packet) { partition.disk.allocCheckLimit(proto.IopsWriteType, 1) if writable := partition.disk.limitWrite.TryRun(currSize, func() { - _, err = store.Write(p.ExtentID, p.ExtentOffset+int64(offset), int64(currSize), data, crc, storage.AppendWriteType, p.IsSyncWrite(), false, false, false) + param := &storage.WriteParam{ + ExtentID: p.ExtentID, + Offset: p.ExtentOffset + int64(offset), + Size: int64(currSize), + Data: data, + Crc: crc, + WriteType: storage.AppendWriteType, + IsSync: p.IsSyncWrite(), + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err = store.Write(param) }); !writable { err = storage.LimitedIoError return diff --git a/master/filecrc.go b/master/filecrc.go index 64cb4455a..bb7c8ff0d 100644 --- a/master/filecrc.go +++ b/master/filecrc.go @@ -89,7 +89,7 @@ func (fc *FileInCore) needCrcRepair(liveReplicas []*DataReplica, getInfoCallback return } if fm.ApplyID == baseApplyId && fm.getFileCrc() != baseCrc { - log.LogErrorf("needCrcRepair. getInfoCallback %v, extent %v, applyID(%v:%v), crc %v", + log.LogWarnf("needCrcRepair. getInfoCallback %v, extent %v, applyID(%v:%v), crc %v", getInfoCallback(), fc.Name, fm.ApplyID, baseApplyId, baseCrc) needRepair = true return diff --git a/storage/extent.go b/storage/extent.go index 19e492ca1..98f140a76 100644 --- a/storage/extent.go +++ b/storage/extent.go @@ -28,6 +28,7 @@ import ( "syscall" "time" + "github.com/cubefs/cubefs/depends/tiglabs/raft/logger" "github.com/cubefs/cubefs/proto" "github.com/cubefs/cubefs/util" "github.com/cubefs/cubefs/util/atomicutil" @@ -45,6 +46,21 @@ const ( ExtentMaxSize = 1024 * 1024 * 1024 * 1024 * 4 // 4TB ) +type WriteParam struct { + ExtentID uint64 + Offset int64 + Size int64 + Data []byte + Crc uint32 + WriteType int + IsSync, IsHole, IsRepair, IsBackupWrite bool +} + +func (wparam *WriteParam) String() (m string) { + return fmt.Sprintf("ExtentID(%v)Offset(%v)Size(%v)Crc(%v)WriteType(%v)IsSync(%v)IsHole(%v)IsRepair(%v)IsBackupWrite(%v)", + wparam.ExtentID, wparam.Offset, wparam.Size, wparam.Crc, wparam.WriteType, wparam.IsSync, wparam.IsHole, wparam.IsRepair, wparam.IsBackupWrite) +} + type ExtentInfo struct { FileID uint64 `json:"fileId"` Size uint64 `json:"size"` @@ -344,28 +360,28 @@ func IsAppendRandomWrite(writeType int) bool { } // WriteTiny performs write on a tiny extent. -func (e *Extent) WriteTiny(data []byte, offset, size int64, crc uint32, writeType int, isSync bool) (err error) { +func (e *Extent) WriteTiny(param *WriteParam) (err error) { e.Lock() defer e.Unlock() - index := offset + size + index := param.Offset + param.Size if index >= ExtentMaxSize { return ExtentIsFullError } - if IsAppendWrite(writeType) && offset != e.dataSize { + if IsAppendWrite(param.WriteType) && param.Offset != e.dataSize { return ParameterMismatchError } - if _, err = e.file.WriteAt(data[:size], int64(offset)); err != nil { + if _, err = e.file.WriteAt(param.Data[:param.Size], int64(param.Offset)); err != nil { return } - if isSync { + if param.IsSync { if err = e.file.Sync(); err != nil { return } } - if !IsAppendWrite(writeType) { + if !IsAppendWrite(param.WriteType) { return } if index%util.PageSize != 0 { @@ -377,85 +393,89 @@ func (e *Extent) WriteTiny(data []byte, offset, size int64, crc uint32, writeTyp } // Write writes data to an extent. -func (e *Extent) Write(data []byte, offset, size int64, crc uint32, writeType int, isSync bool, crcFunc UpdateCrcFunc, ei *ExtentInfo, isHole bool, isRepair bool) (status uint8, err error) { +func (e *Extent) Write(param *WriteParam, crcFunc UpdateCrcFunc) (status uint8, err error) { defer func() { - e.dirty.Store(!isSync) + e.dirty.Store(!param.IsSync) + }() + + if logger.IsEnableDebug() { + log.LogDebugf("action[Extent.Write] path %v write param(%v)", e.filePath, param) + } + defer func() { + e.dirty.Store(!param.IsSync) }() - log.LogDebugf("action[Extent.Write] path %v offset %v size %v writeType %v", e.filePath, offset, size, writeType) status = proto.OpOk if IsTinyExtent(e.extentID) { - err = e.WriteTiny(data, offset, size, crc, writeType, isSync) + err = e.WriteTiny(param) return } - if err = e.checkWriteOffsetAndSize(writeType, offset, size, isRepair); err != nil { - log.LogErrorf("action[Extent.Write] checkWriteOffsetAndSize offset %v size %v writeType %v err %v", - offset, size, writeType, err) - err = newParameterError("extent current size=%d write offset=%d write size=%d", e.dataSize, offset, size) - log.LogInfof("action[Extent.Write] newParameterError path %v offset %v size %v writeType %v err %v", e.filePath, - offset, size, writeType, err) + if err = e.checkWriteOffsetAndSize(param); err != nil { + log.LogErrorf("action[Extent.Write] checkWriteOffsetAndSize write param(%v) err %v", param, err) + err = newParameterError("extent current size=%d write param(%v)", e.dataSize, param) + log.LogInfof("action[Extent.Write] newParameterError path %v write param(%v) err %v", e.filePath, param, err) status = proto.OpTryOtherExtent return } - log.LogDebugf("action[Extent.Write] path %v offset %v size %v writeType %v", e.filePath, offset, size, writeType) // Check if extent file size matches the write offset just in case // multiple clients are writing concurrently. e.Lock() defer e.Unlock() - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v path %v", offset, size, writeType, e.filePath) - if IsAppendWrite(writeType) && e.dataSize != offset { - err = newParameterError("extent current size=%d write offset=%d write size=%d", e.dataSize, offset, size) - log.LogInfof("action[Extent.Write] newParameterError path %v offset %v size %v writeType %v err %v", e.filePath, - offset, size, writeType, err) + + if IsAppendWrite(param.WriteType) && e.dataSize != param.Offset { + err = newParameterError("extent current size=%d write param(%v)", e.dataSize, param) + log.LogInfof("action[Extent.Write] newParameterError path %v write param(%v) err %v", e.filePath, param, err) status = proto.OpTryOtherExtent return } - if IsAppendRandomWrite(writeType) { + if IsAppendRandomWrite(param.WriteType) { if e.snapshotDataOff <= util.ExtentSize { - log.LogInfof("action[Extent.Write] truncate extent %v offset %v size %v writeType %v truncate err %v", e, offset, size, writeType, err) + log.LogInfof("action[Extent.Write] truncate extent %v write param(%v) truncate err %v", e, param, err) if err = e.file.Truncate(util.ExtentSize); err != nil { - log.LogErrorf("action[Extent.Write] offset %v size %v writeType %v truncate err %v", offset, size, writeType, err) + log.LogErrorf("action[Extent.Write] path %v write param(%v) truncate err %v", e.filePath, param, err) return } } } - if isHole { - if err = e.repairPunchHole(offset, size); err != nil { + if param.IsHole { + if err = e.repairPunchHole(param.Offset, param.Size); err != nil { return } } else { - if _, err = e.file.WriteAt(data[:size], int64(offset)); err != nil { - log.LogErrorf("action[Extent.Write] offset %v size %v writeType %v err %v", offset, size, writeType, err) + if _, err = e.file.WriteAt(param.Data[:param.Size], int64(param.Offset)); err != nil { + log.LogErrorf("action[Extent.Write] path %v write param(%v) err %v", e.filePath, param, err) return } } defer func() { - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v path %v", offset, size, writeType, e.filePath) - if IsAppendWrite(writeType) { - atomic.StoreInt64(&e.modifyTime, time.Now().Unix()) - e.dataSize = int64(math.Max(float64(e.dataSize), float64(offset+size))) - log.LogDebugf("action[Extent.Write] e %v offset %v size %v writeType %v", e, offset, size, writeType) - } else if IsAppendRandomWrite(writeType) { - atomic.StoreInt64(&e.modifyTime, time.Now().Unix()) - e.snapshotDataOff = uint64(math.Max(float64(e.snapshotDataOff), float64(offset+size))) + if logger.IsEnableDebug() { + log.LogDebugf("action[Extent.Write] write param(%v),eInfo %v,err %v, status %v", param, e, err, status) + } + if IsAppendWrite(param.WriteType) { + atomic.StoreInt64(&e.modifyTime, time.Now().Unix()) + e.dataSize = int64(math.Max(float64(e.dataSize), float64(param.Offset+param.Size))) + log.LogDebugf("action[Extent.Write] eInfo %v write param(%v)", e, param) + } else if IsAppendRandomWrite(param.WriteType) { + atomic.StoreInt64(&e.modifyTime, time.Now().Unix()) + e.snapshotDataOff = uint64(math.Max(float64(e.snapshotDataOff), float64(param.Offset+param.Size))) + } + if logger.IsEnableDebug() { + log.LogDebugf("action[Extent.Write] write param(%v) dataSize %v snapshotDataOff %v", param, e.dataSize, e.snapshotDataOff) } - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v dataSize %v snapshotDataOff %v", - offset, size, writeType, e.dataSize, e.snapshotDataOff) }() - if isSync { + if param.IsSync { if err = e.file.Sync(); err != nil { - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v err %v", - offset, size, writeType, err) + log.LogDebugf("action[Extent.Write] write param(%v) err %v", param, err) return } } // NOTE: compute crc - beginOffset := offset - endOffset := offset + size + beginOffset := param.Offset + endOffset := param.Offset + param.Size for beginOffset != endOffset { // NOTE: take a block blockNo := beginOffset / util.BlockSize @@ -469,14 +489,14 @@ func (e *Extent) Write(data []byte, offset, size int64, crc uint32, writeType in // NOTE: aliagn, compute crc if offsetInBlock == 0 && sizeInBlock == util.BlockSize { - err = crcFunc(e, int(blockNo), crc) - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v err %v crcOffset %v", offset, size, writeType, err, beginOffset) + err = crcFunc(e, int(blockNo), param.Crc) + log.LogDebugf("action[Extent.Write] write param(%v) err %v crcOffset %v", param, err, beginOffset) beginOffset += sizeInBlock continue } // NOTE: not aliagn err = crcFunc(e, int(blockNo), 0) - log.LogDebugf("action[Extent.Write] offset %v size %v writeType %v err %v crcOffset %v", offset, size, writeType, err, beginOffset) + log.LogDebugf("action[Extent.Write] write param(%v) err %v crcOffset %v", param, err, beginOffset) beginOffset += sizeInBlock } return @@ -521,20 +541,20 @@ func (e *Extent) checkReadOffsetAndSize(offset, size int64) error { return nil } -func (e *Extent) checkWriteOffsetAndSize(writeType int, offset, size int64, isRepair bool) error { - err := newParameterError("writeType=%d offset=%d size=%d", writeType, offset, size) - if IsAppendWrite(writeType) { - if size == 0 || - offset+size > util.ExtentSize || - offset >= util.ExtentSize { +func (e *Extent) checkWriteOffsetAndSize(param *WriteParam) error { + err := newParameterError("writeType=%d offset=%d size=%d", param.WriteType, param.Offset, param.Size) + if IsAppendWrite(param.WriteType) { + if param.Size == 0 || + param.Offset+param.Size > util.ExtentSize || + param.Offset >= util.ExtentSize { return err } - if !isRepair && size > util.BlockSize { + if !param.IsRepair && param.Size > util.BlockSize { return err } - } else if IsAppendRandomWrite(writeType) { - log.LogDebugf("action[checkOffsetAndSize] offset %v size %v", offset, size) - if offset < util.ExtentSize || size == 0 { + } else if IsAppendRandomWrite(param.WriteType) { + log.LogDebugf("action[checkOffsetAndSize] offset %v size %v", param.Offset, param.Size) + if param.Offset < util.ExtentSize || param.Size == 0 { return err } } diff --git a/storage/extent_store.go b/storage/extent_store.go index e32026d03..9d35c5d91 100644 --- a/storage/extent_store.go +++ b/storage/extent_store.go @@ -149,9 +149,8 @@ type ExtentStore struct { extentLock bool stopMutex sync.RWMutex stopC chan interface{} - - ApplyId uint64 - ApplyIdMutex sync.RWMutex + ApplyIdMutex sync.RWMutex + ApplyId uint64 } func MkdirAll(name string) (err error) { @@ -288,7 +287,9 @@ func (s *ExtentStore) SnapShot() (files []*proto.File, err error) { for _, ei := range normalExtentSnapshot { file := GetSnapShotFileFromPool() file.Name = strconv.FormatUint(ei.FileID, 10) + file.Size = uint32(ei.Size) + file.Modified = ei.ModifyTime file.Crc = atomic.LoadUint32(&ei.Crc) file.ApplyID = ei.ApplyID @@ -626,12 +627,12 @@ func (s *ExtentStore) initBaseFileID() error { } // Write writes the given extent to the disk. -func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, crc uint32, writeType int, isSync bool, isHole bool, isRepair, isBackupWrite bool) (status uint8, err error) { +func (s *ExtentStore) Write(param *WriteParam) (status uint8, err error) { s.stopMutex.RLock() defer s.stopMutex.RUnlock() if s.IsClosed() { err = ErrStoreAlreadyClosed - log.LogErrorf("[Write] store(%v) failed to write extent(%v), err(%v)", s.dataPath, extentID, err) + log.LogErrorf("[Write] store(%v) failed to write param(%v), err(%v)", s.dataPath, param, err) return } @@ -641,22 +642,22 @@ func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, cr ) s.elMutex.RLock() - if isBackupWrite { + if param.IsBackupWrite { // NOTE: meet an error is impossible - _, ok := s.GetExtentInfo(extentID) + _, ok := s.GetExtentInfo(param.ExtentID) if !ok { s.elMutex.RUnlock() - err = fmt.Errorf("extent(%v) is not locked", extentID) - log.LogErrorf("[Write] gc_extent[%d] is not locked", extentID) + err = fmt.Errorf("extent(%v) is not locked", param.ExtentID) + log.LogErrorf("[Write] gc_extent[%d] is not locked", param.ExtentID) return } } else { if s.extentLock { - if flag, ok := s.extentLockMap[extentID]; ok { - log.LogErrorf("[Write] gc_extent_lock[%d] is locked, path %s", extentID, s.dataPath) + if flag, ok := s.extentLockMap[param.ExtentID]; ok { + log.LogErrorf("[Write] gc_extent_lock[%d] is locked, path %s", param.ExtentID, s.dataPath) if flag == proto.GcDeleteFlag { s.elMutex.RUnlock() - err = fmt.Errorf("extent(%v) is locked", extentID) + err = fmt.Errorf("extent(%v) is locked", param.ExtentID) return } } @@ -666,7 +667,7 @@ func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, cr s.eiMutex.Lock() status = proto.OpOk - ei = s.extentInfoMap[extentID] + ei = s.extentInfoMap[param.ExtentID] e, err = s.extentWithHeader(ei) s.eiMutex.Unlock() if err != nil { @@ -674,13 +675,13 @@ func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, cr } // update access time atomic.StoreInt64(&ei.AccessTime, time.Now().Unix()) - log.LogDebugf("action[Write] dp %v extentID %v offset %v size %v writeTYPE %v isRepair(%v)", s.partitionID, extentID, offset, size, writeType, isRepair) - if err = s.checkOffsetAndSize(extentID, offset, size, writeType, isRepair); err != nil { + log.LogDebugf("action[Write] dp %v write param(%v)", s.partitionID, param) + if err = s.checkOffsetAndSize(param); err != nil { log.LogInfof("action[Write] path %v err %v", e.filePath, err) return status, err } - status, err = e.Write(data, offset, size, crc, writeType, isSync, s.PersistenceBlockCrc, ei, isHole, isRepair) + status, err = e.Write(param, s.PersistenceBlockCrc) if err != nil { log.LogInfof("action[Write] path %v err %v", e.filePath, err) return status, err @@ -690,27 +691,27 @@ func (s *ExtentStore) Write(extentID uint64, offset, size int64, data []byte, cr return status, nil } -func (s *ExtentStore) checkOffsetAndSize(extentID uint64, offset, size int64, writeType int, isRepair bool) error { - if IsTinyExtent(extentID) { +func (s *ExtentStore) checkOffsetAndSize(param *WriteParam) error { + if IsTinyExtent(param.ExtentID) { return nil } // random write pos can happen on modAppend partition of extent - if writeType == RandomWriteType { + if param.WriteType == RandomWriteType { return nil } - if writeType == AppendRandomWriteType { - if offset < util.ExtentSize { - return newParameterError("writeType=%d offset=%d size=%d", writeType, offset, size) + if param.WriteType == AppendRandomWriteType { + if param.Offset < util.ExtentSize { + return newParameterError("Write param error(%v)", param) } return nil } - if size == 0 || - offset >= util.BlockCount*util.BlockSize || - offset+size > util.BlockCount*util.BlockSize { - return newParameterError("offset=%d size=%d", offset, size) + if param.Size == 0 || + param.Offset >= util.BlockCount*util.BlockSize || + param.Offset+param.Size > util.BlockCount*util.BlockSize { + return newParameterError("offset=%d size=%d", param.Offset, param.Size) } - if !isRepair && size > util.BlockSize { - return newParameterError("offset=%d size=%d", offset, size) + if !param.IsRepair && param.Size > util.BlockSize { + return newParameterError("offset=%d size=%d", param.Offset, param.Size) } return nil } diff --git a/storage/extent_store_test.go b/storage/extent_store_test.go index 0f44683d0..a59df6770 100644 --- a/storage/extent_store_test.go +++ b/storage/extent_store_test.go @@ -40,13 +40,27 @@ func extentStoreNormalRwTest(t *testing.T, s *storage.ExtentStore, id uint64) { data := []byte(dataStr) crc := crc32.ChecksumIEEE(data) // append write - _, err := s.Write(id, 0, int64(len(data)), data, crc, storage.AppendWriteType, true, false, false, false) + param := &storage.WriteParam{ + ExtentID: id, + Offset: 0, + Size: int64(len(data)), + Data: data, + Crc: crc, + WriteType: storage.AppendWriteType, + IsSync: true, + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + + _, err := s.Write(param) require.NoError(t, err) actualCrc, err := s.Read(id, 0, int64(len(data)), data, false, false) require.NoError(t, err) require.EqualValues(t, crc, actualCrc) // random write - _, err = s.Write(id, 0, int64(len(data)), data, crc, storage.RandomWriteType, true, false, false, false) + param.WriteType = storage.RandomWriteType + _, err = s.Write(param) require.NoError(t, err) actualCrc, err = s.Read(id, 0, int64(len(data)), data, false, false) require.NoError(t, err) @@ -93,7 +107,19 @@ func extentMarkDeleteTinyTest(t *testing.T, s *storage.ExtentStore, id uint64) { require.NotEqualValues(t, size, 0) // write second file to extent crc := crc32.ChecksumIEEE(data) - _, err = s.Write(id, size, int64(len(data)), data, crc, storage.AppendWriteType, true, false, false, false) + param := &storage.WriteParam{ + ExtentID: id, + Offset: size, + Size: int64(len(data)), + Data: data, + Crc: crc, + WriteType: storage.AppendWriteType, + IsSync: true, + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err = s.Write(param) require.NoError(t, err) // mark delete first file extentStoreMarkDeleteTiny(t, s, id, 0, size) @@ -155,7 +181,19 @@ func reopenExtentStoreTest(t *testing.T, dpType int) { data := []byte(dataStr) crc := crc32.ChecksumIEEE(data) // write some data - _, err = s.Write(id, 0, int64(len(data)), data, crc, storage.AppendWriteType, true, false, false, false) + param := &storage.WriteParam{ + ExtentID: id, + Offset: 0, + Size: int64(len(data)), + Data: data, + Crc: crc, + WriteType: storage.AppendWriteType, + IsSync: true, + IsHole: false, + IsRepair: false, + IsBackupWrite: false, + } + _, err = s.Write(param) require.NoError(t, err) firstSnap, err := s.SnapShot() require.NoError(t, err) diff --git a/storage/extent_test.go b/storage/extent_test.go index c4753e8e8..c62e66efc 100644 --- a/storage/extent_test.go +++ b/storage/extent_test.go @@ -56,51 +56,82 @@ func getMockCrcPersist(t *testing.T) storage.UpdateCrcFunc { func normalExtentRwTest(t *testing.T, e *storage.Extent) { data := []byte(dataStr) - _, err := e.Write(data, 0, 0, 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + param := &storage.WriteParam{ + Data: data, + Offset: 0, + ExtentID: 0, + Crc: 0, + WriteType: storage.AppendWriteType, + IsSync: true, + IsHole: false, + IsRepair: false, + } + _, err := e.Write(param, getMockCrcPersist(t)) + param.Size = int64(len(data)) require.Error(t, err) // append write - _, err = e.Write(data, 0, int64(len(data)), 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) require.EqualValues(t, e.Size(), len(data)) _, err = e.Read(data, 0, int64(len(data)), false) require.NoError(t, err) require.Equal(t, string(data), dataStr) // failed append write - _, err = e.Write(data, 0, int64(len(data)), 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + param.WriteType = storage.AppendWriteType + _, err = e.Write(param, getMockCrcPersist(t)) require.Error(t, err) // random append write oldSize := e.Size() - _, err = e.Write(data, 0, int64(len(data)), 0, storage.RandomWriteType, true, getMockCrcPersist(t), nil, false, false) + param.WriteType = storage.RandomWriteType + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) require.Equal(t, e.Size(), oldSize) _, err = e.Read(data, 0, int64(len(data)), false) require.NoError(t, err) require.Equal(t, string(data), dataStr) - _, err = e.Write(data, util.BlockSize, dataSize, 0, storage.RandomWriteType, true, getMockCrcPersist(t), nil, false, false) + + param.Offset = util.BlockSize + param.Size = dataSize + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) - _, err = e.Write(data, util.ExtentSize, dataSize, 0, storage.RandomWriteType, true, getMockCrcPersist(t), nil, false, false) + param.Offset = util.ExtentSize + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) // TODO: append random write test } func tinyExtentRwTest(t *testing.T, e *storage.Extent) { data := []byte(dataStr) + param := &storage.WriteParam{ + Data: data, + IsSync: true, + } // write oversize - _, err := e.Write(data, storage.ExtentMaxSize, dataSize, 0, storage.RandomWriteType, true, getMockCrcPersist(t), nil, false, false) + param.Offset = storage.ExtentMaxSize + param.Size = dataSize + param.WriteType = storage.RandomWriteType + _, err := e.Write(param, getMockCrcPersist(t)) require.ErrorIs(t, err, storage.ExtentIsFullError) + // append write - _, err = e.Write(data, 0, int64(len(data)), 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + param.Offset = 0 + param.Size = int64(len(data)) + param.WriteType = storage.AppendWriteType + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) require.EqualValues(t, e.Size()%util.PageSize, 0) _, err = e.Read(data, 0, int64(len(data)), false) require.NoError(t, err) require.Equal(t, string(data), dataStr) + // failed append write - _, err = e.Write(data, 0, int64(len(data)), 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + _, err = e.Write(param, getMockCrcPersist(t)) require.Error(t, err) // random write oldSize := e.Size() - _, err = e.Write(data, int64(len(data)), int64(len(data)), 0, storage.RandomWriteType, true, getMockCrcPersist(t), nil, false, false) + param.WriteType = storage.RandomWriteType + param.Offset = int64(len(data)) + _, err = e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) require.Equal(t, e.Size(), oldSize) _, err = e.Read(data, int64(len(data)), int64(len(data)), false) @@ -285,12 +316,23 @@ func TestExtentRecovery(t *testing.T) { headSize := 128 * 1024 data := bytes.Repeat([]byte("s"), headSize) + param := &storage.WriteParam{ + Data: data, + IsSync: true, + } for i := 0; i < 10; i++ { - _, err := e.Write(data, int64(i)*util.BlockSize, int64(headSize), 0, storage.AppendWriteType, true, getMockCrcPersist(t), nil, false, false) + param.Offset = int64(i) * util.BlockSize + param.Size = int64(headSize) + param.WriteType = storage.AppendWriteType + + _, err := e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) } for i := 0; i < 10; i++ { - _, err := e.Write(data, int64(i)*util.BlockSize+util.ExtentSize, int64(headSize), 0, storage.AppendRandomWriteType, true, getMockCrcPersist(t), nil, false, false) + param.Offset = int64(i)*util.BlockSize + util.ExtentSize + param.Size = int64(headSize) + param.WriteType = storage.AppendRandomWriteType + _, err := e.Write(param, getMockCrcPersist(t)) require.NoError(t, err) } e.GetFile().Close()