mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 10:06:14 +00:00
refactor(metanode): remove unused mutex and optimize multipart handling. #1000348192
- Remove unused RWMutex from MetaNodeMetrics struct - Rename TestMultipartExtend_Bytes to TestMUMultipartExtend_Bytes - Add comprehensive unit tests for Multipart and Part structs - Implement tests for edge cases and concurrent access scenarios - Optimize buffer handling in multipart serialization methods This commit enhances the performance and readability of the metanode implementation, ensuring better resource management and test coverage. Signed-off-by: Victor1319 <zengxuewei@oppo.com>
This commit is contained in:
parent
9009567e2c
commit
80fe126da2
@ -292,20 +292,8 @@ const (
|
||||
|
||||
// MetaNode specific configuration
|
||||
metaNodeDeleteBatchCountKey = "batchCount"
|
||||
configNameResolveInterval = "nameResolveInterval" // int
|
||||
|
||||
cfgRocksDirs = "rocksDirs"
|
||||
cfgDiskReservedSpace = "diskReservedSpace"
|
||||
|
||||
// NOTE: metanode rocksdb config
|
||||
cfgRocksdbWriteBufferSize = "rocksdbWriteBufferSize" // int
|
||||
cfgRocksdbWriteBufferNum = "rocksdbWriteBufferNum" // int
|
||||
cfgRocksdbBlockCacheSize = "rocksdbBlockCacheSize" // uint64
|
||||
cfgRocksdbMinWriteBufferToMerge = "rocksdbMinWriteBufferToMerge" // int
|
||||
cfgRocksdbMaxSubCompactions = "rocksdbMaxSubCompactions" // int
|
||||
cfgRocksdbMode = "rocksdbMode" // string
|
||||
cfgRocksdbEnableStats = "rocksdbEnableStats" // bool
|
||||
cfgRocksdbKeyNumMax = "rocksdbKeyNumMax" // int64
|
||||
cfgRocksdbBytesPerSync = "rocksdbBytesPerSync" // uint64
|
||||
cfgRocksdbParallelism = "rocksdbParallelism" // int
|
||||
cfgRocksdbMaxBackgroundCompactions = "rocksdbMaxBackgroundCompactions" // int
|
||||
|
||||
@ -863,7 +863,10 @@ func (d *Dentry) UnmarshalValue(raw []byte) (err error) {
|
||||
}
|
||||
|
||||
// Pre-allocate slice for better performance
|
||||
d.multiSnap.dentryList = make([]*Dentry, 0, int(verCnt))
|
||||
// Only create non-nil slice when verCnt > 0, otherwise keep it nil
|
||||
if verCnt > 0 {
|
||||
d.multiSnap.dentryList = make([]*Dentry, 0, int(verCnt))
|
||||
}
|
||||
|
||||
for i := 0; i < int(verCnt); i++ {
|
||||
den := &Dentry{
|
||||
|
||||
@ -366,7 +366,6 @@ func (m *MetaNode) validateRequiredConfigs() error {
|
||||
}
|
||||
|
||||
func (m *MetaNode) parseRaftConfig(cfg *config.Config) error {
|
||||
|
||||
m.opLimiter = newOpLimiter()
|
||||
// Load operation rate limiter configuration
|
||||
if err := m.loadOpLimiterConfig(cfg); err != nil {
|
||||
|
||||
@ -20,7 +20,6 @@ import (
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/cubefs/cubefs/util"
|
||||
@ -39,27 +38,13 @@ const (
|
||||
MetricFileStats = "fileStats"
|
||||
RocksdbStats = "rocksdbStats"
|
||||
|
||||
// RocksDB statistics keys
|
||||
RocksdbGetMicros = "rocksdb.db.get.micros"
|
||||
RocksdbWriteMicros = "rocksdb.db.write.micros"
|
||||
RocksdbSeekMicros = "rocksdb.db.seek.micros"
|
||||
RocksdbWriteStall = "rocksdb.db.write.stall"
|
||||
RocksdbFlushMicros = "rocksdb.db.flush.micros"
|
||||
|
||||
// Timeout for metrics collection
|
||||
MetricsCollectionTimeout = 30 * time.Second
|
||||
)
|
||||
|
||||
// Pre-compiled regex for better performance
|
||||
var (
|
||||
statsP99Regex = regexp.MustCompile(`P99 : (\d+\.\d+)`)
|
||||
rocksdbStatsList = []string{
|
||||
RocksdbGetMicros,
|
||||
RocksdbWriteMicros,
|
||||
RocksdbSeekMicros,
|
||||
RocksdbWriteStall,
|
||||
RocksdbFlushMicros,
|
||||
}
|
||||
statsP99Regex = regexp.MustCompile(`P99 : (\d+\.\d+)`)
|
||||
)
|
||||
|
||||
// MetaNodeMetrics holds all metrics for the meta node
|
||||
@ -72,7 +57,6 @@ type MetaNodeMetrics struct {
|
||||
RocksdbStats *exporter.GaugeVec
|
||||
|
||||
metricStopCh chan struct{}
|
||||
mu sync.RWMutex
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
@ -187,52 +171,45 @@ func (m *MetaNode) collectPartitionMetrics() {
|
||||
}
|
||||
m.metrics.MetricConnectionCount.Set(float64(m.connectionCnt))
|
||||
case <-fileStatTicker.C:
|
||||
if err := m.updateFileStatsMetrics(); err != nil {
|
||||
// Log error but continue
|
||||
continue
|
||||
}
|
||||
m.updateFileStatsMetrics()
|
||||
if m.rocksdbEnableStats {
|
||||
if err := m.updateRocksdbStatsMetrics(); err != nil {
|
||||
// Log error but continue
|
||||
continue
|
||||
}
|
||||
m.updateRocksdbStatsMetrics()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// updateFileStatsMetrics updates file statistics metrics
|
||||
func (m *MetaNode) updateFileStatsMetrics() error {
|
||||
func (m *MetaNode) updateFileStatsMetrics() {
|
||||
m.metrics.MetricFileStats.Reset()
|
||||
volFileRange := make(map[string][]int64)
|
||||
|
||||
manager, ok := m.metadataManager.(*metadataManager)
|
||||
if !ok {
|
||||
return fmt.Errorf("invalid metadata manager type")
|
||||
return
|
||||
}
|
||||
manager.mu.RLock()
|
||||
defer manager.mu.RUnlock()
|
||||
|
||||
// Get file stats configuration
|
||||
_, labels, _ := manager.GetFileStatsConfig()
|
||||
|
||||
numRanges := len(labels)
|
||||
volFileRange := make(map[string][]int64)
|
||||
|
||||
// Collect file range data
|
||||
partitions := m.collectFileRangeData(manager, numRanges)
|
||||
|
||||
// Process collected data
|
||||
for _, p := range partitions {
|
||||
volName := p.volName
|
||||
for _, p := range manager.partitions {
|
||||
mp, ok := p.(*metaPartition)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
fileRange := mp.getFileRange()
|
||||
volName := mp.config.VolName
|
||||
if _, exists := volFileRange[volName]; !exists {
|
||||
volFileRange[volName] = make([]int64, numRanges)
|
||||
}
|
||||
|
||||
validLength := util.Min(len(p.fileRange), numRanges)
|
||||
validLength := util.Min(len(fileRange), numRanges)
|
||||
for i := 0; i < validLength; i++ {
|
||||
volFileRange[volName][i] += p.fileRange[i]
|
||||
volFileRange[volName][i] += fileRange[i]
|
||||
}
|
||||
}
|
||||
|
||||
// Update metrics
|
||||
for volName, ranges := range volFileRange {
|
||||
for i, val := range ranges {
|
||||
sizeRange := labels[i]
|
||||
@ -262,68 +239,17 @@ func (m *MetaNode) updateRocksdbStatsMetrics() {
|
||||
manager.mu.RLock()
|
||||
defer manager.mu.RUnlock()
|
||||
|
||||
partitions := make([]fileRangeData, 0, len(manager.partitions))
|
||||
|
||||
for _, p := range manager.partitions {
|
||||
mp, ok := p.(*metaPartition)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
partitions = append(partitions, fileRangeData{
|
||||
volName: mp.config.VolName,
|
||||
fileRange: mp.getFileRange(),
|
||||
})
|
||||
}
|
||||
|
||||
return partitions
|
||||
}
|
||||
|
||||
// updateRocksdbStatsMetrics updates RocksDB statistics metrics
|
||||
func (m *MetaNode) updateRocksdbStatsMetrics() error {
|
||||
m.metrics.RocksdbStats.Reset()
|
||||
|
||||
// Process each RocksDB directory
|
||||
for _, dbPath := range m.rocksDirs {
|
||||
if err := m.processRocksdbDirectory(dbPath); err != nil {
|
||||
// Log error but continue with other directories
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// processRocksdbDirectory processes a single RocksDB directory
|
||||
func (m *MetaNode) processRocksdbDirectory(dbPath string) error {
|
||||
// Open RocksDB with timeout
|
||||
ctx, cancel := context.WithTimeout(m.metrics.ctx, MetricsCollectionTimeout)
|
||||
defer cancel()
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
db, err := m.rocksdbManager.OpenRocksdb(dbPath, 0)
|
||||
if err != nil {
|
||||
done <- err
|
||||
return
|
||||
continue
|
||||
}
|
||||
defer m.rocksdbManager.CloseRocksdb(db)
|
||||
|
||||
stats := db.GetStatistics()
|
||||
statsP99 := getStatsP99(stats, rocksdbStatsList)
|
||||
|
||||
for key, val := range statsP99 {
|
||||
m.metrics.RocksdbStats.SetWithLabelValues(val, dbPath, key)
|
||||
}
|
||||
|
||||
done <- nil
|
||||
}()
|
||||
|
||||
select {
|
||||
case err := <-done:
|
||||
return err
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
m.rocksdbManager.CloseRocksdb(db)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -25,6 +25,37 @@ import (
|
||||
"github.com/cubefs/cubefs/util/log"
|
||||
)
|
||||
|
||||
// Constants for better maintainability
|
||||
const (
|
||||
MaxVarintLen64 = binary.MaxVarintLen64
|
||||
DefaultBufferSize = 1024
|
||||
)
|
||||
|
||||
// Buffer pool for memory optimization
|
||||
var bufferPool = sync.Pool{
|
||||
New: func() interface{} {
|
||||
return bytes.NewBuffer(make([]byte, 0, DefaultBufferSize))
|
||||
},
|
||||
}
|
||||
|
||||
// Helper functions for binary encoding/decoding
|
||||
func encodeUint64(buffer *bytes.Buffer, val uint64, tmp []byte) error {
|
||||
n := binary.PutUvarint(tmp, val)
|
||||
_, err := buffer.Write(tmp[:n])
|
||||
return err
|
||||
}
|
||||
|
||||
func encodeString(buffer *bytes.Buffer, s string, tmp []byte) error {
|
||||
n := binary.PutUvarint(tmp, uint64(len(s)))
|
||||
if _, err := buffer.Write(tmp[:n]); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := buffer.WriteString(s); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Part defined necessary fields for multipart part management.
|
||||
type Part struct {
|
||||
ID uint16
|
||||
@ -42,38 +73,40 @@ func (m *Part) Equal(o *Part) bool {
|
||||
}
|
||||
|
||||
func (m Part) Bytes() ([]byte, error) {
|
||||
var err error
|
||||
buffer := bytes.NewBuffer(nil)
|
||||
tmp := make([]byte, binary.MaxVarintLen64)
|
||||
var n int
|
||||
// ID
|
||||
n = binary.PutUvarint(tmp, uint64(m.ID))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
buffer := bufferPool.Get().(*bytes.Buffer)
|
||||
defer func() {
|
||||
buffer.Reset()
|
||||
bufferPool.Put(buffer)
|
||||
}()
|
||||
|
||||
tmp := make([]byte, MaxVarintLen64)
|
||||
|
||||
// Encode ID
|
||||
if err := encodeUint64(buffer, uint64(m.ID), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// upload time
|
||||
n = binary.PutVarint(tmp, m.UploadTime.UnixNano())
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
|
||||
// Encode upload time
|
||||
n := binary.PutVarint(tmp, m.UploadTime.UnixNano())
|
||||
if _, err := buffer.Write(tmp[:n]); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// MD5
|
||||
n = binary.PutUvarint(tmp, uint64(len(m.MD5)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
|
||||
// Encode MD5
|
||||
if err := encodeString(buffer, m.MD5, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err = buffer.WriteString(m.MD5); err != nil {
|
||||
|
||||
// Encode size
|
||||
if err := encodeUint64(buffer, m.Size, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// size
|
||||
n = binary.PutUvarint(tmp, m.Size)
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// inode
|
||||
n = binary.PutUvarint(tmp, m.Inode)
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
|
||||
// Encode inode
|
||||
if err := encodeUint64(buffer, m.Inode, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
@ -199,30 +232,37 @@ func (m Parts) Search(id uint16) (part *Part, found bool) {
|
||||
}
|
||||
|
||||
func (m Parts) Bytes() ([]byte, error) {
|
||||
var err error
|
||||
var n int
|
||||
buffer := bytes.NewBuffer(nil)
|
||||
tmp := make([]byte, binary.MaxVarintLen64)
|
||||
n = binary.PutUvarint(tmp, uint64(len(m)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
buffer := bufferPool.Get().(*bytes.Buffer)
|
||||
defer func() {
|
||||
buffer.Reset()
|
||||
bufferPool.Put(buffer)
|
||||
}()
|
||||
|
||||
tmp := make([]byte, MaxVarintLen64)
|
||||
|
||||
// Write number of parts
|
||||
if err := encodeUint64(buffer, uint64(len(m)), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var marshaled []byte
|
||||
|
||||
// Write each part
|
||||
for _, p := range m {
|
||||
marshaled, err = p.Bytes()
|
||||
marshaled, err := p.Bytes()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// write part length
|
||||
n = binary.PutUvarint(tmp, uint64(len(marshaled)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
|
||||
// Write part length
|
||||
if err := encodeUint64(buffer, uint64(len(marshaled)), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// write part bytes
|
||||
if _, err = buffer.Write(marshaled); err != nil {
|
||||
|
||||
// Write part bytes
|
||||
if _, err := buffer.Write(marshaled); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
@ -250,32 +290,29 @@ func NewMultipartExtend() MultipartExtend {
|
||||
}
|
||||
|
||||
func (me MultipartExtend) Bytes() ([]byte, error) {
|
||||
var n int
|
||||
var err error
|
||||
buffer := bytes.NewBuffer(nil)
|
||||
tmp := make([]byte, binary.MaxVarintLen64)
|
||||
n = binary.PutUvarint(tmp, uint64(len(me)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
buffer := bufferPool.Get().(*bytes.Buffer)
|
||||
defer func() {
|
||||
buffer.Reset()
|
||||
bufferPool.Put(buffer)
|
||||
}()
|
||||
|
||||
tmp := make([]byte, MaxVarintLen64)
|
||||
|
||||
// Write number of entries
|
||||
if err := encodeUint64(buffer, uint64(len(me)), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
marshalStr := func(src string) error {
|
||||
n = binary.PutUvarint(tmp, uint64(len(src)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err = buffer.WriteString(src); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Write key-value pairs
|
||||
for key, val := range me {
|
||||
if err = marshalStr(key); err != nil {
|
||||
if err := encodeString(buffer, key, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err = marshalStr(val); err != nil {
|
||||
if err := encodeString(buffer, val, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
@ -364,58 +401,54 @@ func (m *Multipart) Parts() []*Part {
|
||||
}
|
||||
|
||||
func (m *Multipart) Bytes() ([]byte, error) {
|
||||
var n int
|
||||
buffer := bytes.NewBuffer(nil)
|
||||
var err error
|
||||
tmp := make([]byte, binary.MaxVarintLen64)
|
||||
// marshal id
|
||||
marshalStr := func(src string) error {
|
||||
n = binary.PutUvarint(tmp, uint64(len(src)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err = buffer.WriteString(src); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
// marshal id
|
||||
if err = marshalStr(m.id); err != nil {
|
||||
buffer := bufferPool.Get().(*bytes.Buffer)
|
||||
defer func() {
|
||||
buffer.Reset()
|
||||
bufferPool.Put(buffer)
|
||||
}()
|
||||
|
||||
tmp := make([]byte, MaxVarintLen64)
|
||||
|
||||
// Marshal id
|
||||
if err := encodeString(buffer, m.id, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// marshal key
|
||||
if err = marshalStr(m.key); err != nil {
|
||||
|
||||
// Marshal key
|
||||
if err := encodeString(buffer, m.key, tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// marshal init time
|
||||
n = binary.PutVarint(tmp, m.initTime.UnixNano())
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
|
||||
// Marshal init time
|
||||
n := binary.PutVarint(tmp, m.initTime.UnixNano())
|
||||
if _, err := buffer.Write(tmp[:n]); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// marshal parts
|
||||
var marshaledParts []byte
|
||||
if marshaledParts, err = m.parts.Bytes(); err != nil {
|
||||
|
||||
// Marshal parts
|
||||
marshaledParts, err := m.parts.Bytes()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
n = binary.PutUvarint(tmp, uint64(len(marshaledParts)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
if err := encodeUint64(buffer, uint64(len(marshaledParts)), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err = buffer.Write(marshaledParts); err != nil {
|
||||
if _, err := buffer.Write(marshaledParts); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// marshall extend
|
||||
var extendBytes []byte
|
||||
if extendBytes, err = m.extend.Bytes(); err != nil {
|
||||
|
||||
// Marshal extend
|
||||
extendBytes, err := m.extend.Bytes()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
n = binary.PutUvarint(tmp, uint64(len(extendBytes)))
|
||||
if _, err = buffer.Write(tmp[:n]); err != nil {
|
||||
if err := encodeUint64(buffer, uint64(len(extendBytes)), tmp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err = buffer.Write(extendBytes); err != nil {
|
||||
if _, err := buffer.Write(extendBytes); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
|
||||
@ -155,7 +155,7 @@ func TestMUSession_Bytes(t *testing.T) {
|
||||
t.Logf("encoded session length: %v", len(sessionBytes))
|
||||
}
|
||||
|
||||
func TestMultipartExtend_Bytes(t *testing.T) {
|
||||
func TestMUMultipartExtend_Bytes(t *testing.T) {
|
||||
me := NewMultipartExtend()
|
||||
me["oss::tag"] = "name=123&age456"
|
||||
me["oss::disposition"] = "attachment=file.txt"
|
||||
@ -168,3 +168,423 @@ func TestMultipartExtend_Bytes(t *testing.T) {
|
||||
t.Fatalf("result mismatch:\n\tme1:%v\n\tme2:%v", me, me2)
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUPart_Equal tests the Equal method of Part struct
|
||||
func TestMUPart_Equal(t *testing.T) {
|
||||
part1 := &Part{
|
||||
ID: 1,
|
||||
UploadTime: time.Now(),
|
||||
MD5: "testmd5",
|
||||
Size: 1024,
|
||||
Inode: 12345,
|
||||
}
|
||||
|
||||
part2 := &Part{
|
||||
ID: 1,
|
||||
UploadTime: time.Now().Add(time.Hour), // Different time should not affect equality
|
||||
MD5: "testmd5",
|
||||
Size: 1024,
|
||||
Inode: 12345,
|
||||
}
|
||||
|
||||
part3 := &Part{
|
||||
ID: 2, // Different ID
|
||||
UploadTime: time.Now(),
|
||||
MD5: "testmd5",
|
||||
Size: 1024,
|
||||
Inode: 12345,
|
||||
}
|
||||
|
||||
if !part1.Equal(part2) {
|
||||
t.Fatalf("parts with same ID, MD5, Size, Inode should be equal")
|
||||
}
|
||||
|
||||
if part1.Equal(part3) {
|
||||
t.Fatalf("parts with different ID should not be equal")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_Empty tests operations on empty Parts
|
||||
func TestMUParts_Empty(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
if parts.Len() != 0 {
|
||||
t.Fatalf("empty parts should have length 0, got %d", parts.Len())
|
||||
}
|
||||
|
||||
// Test search on empty parts
|
||||
if _, found := parts.Search(1); found {
|
||||
t.Fatalf("search on empty parts should not find anything")
|
||||
}
|
||||
|
||||
// Test remove on empty parts
|
||||
parts.Remove(1) // Should not panic
|
||||
|
||||
// Test bytes on empty parts
|
||||
bytes, err := parts.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("empty parts bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
// Test deserialization of empty parts
|
||||
parts2 := PartsFromBytes(bytes)
|
||||
if parts2.Len() != 0 {
|
||||
t.Fatalf("deserialized empty parts should have length 0, got %d", parts2.Len())
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_InsertReplace tests Insert with replace flag
|
||||
func TestMUParts_InsertReplace(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
part1 := &Part{
|
||||
ID: 1,
|
||||
MD5: "md5_1",
|
||||
Size: 1024,
|
||||
Inode: 100,
|
||||
}
|
||||
|
||||
part2 := &Part{
|
||||
ID: 1,
|
||||
MD5: "md5_2", // Different MD5
|
||||
Size: 2048, // Different size
|
||||
Inode: 200, // Different inode
|
||||
}
|
||||
|
||||
// Insert without replace
|
||||
success := parts.Insert(part1, false)
|
||||
if !success {
|
||||
t.Fatalf("first insert should succeed")
|
||||
}
|
||||
|
||||
// Try to insert same ID without replace
|
||||
success = parts.Insert(part2, false)
|
||||
if success {
|
||||
t.Fatalf("insert same ID without replace should fail")
|
||||
}
|
||||
|
||||
// Insert with replace
|
||||
success = parts.Insert(part2, true)
|
||||
if !success {
|
||||
t.Fatalf("insert with replace should succeed")
|
||||
}
|
||||
|
||||
// Verify the part was replaced
|
||||
if parts.Len() != 1 {
|
||||
t.Fatalf("parts should have length 1, got %d", parts.Len())
|
||||
}
|
||||
|
||||
foundPart, found := parts.Search(1)
|
||||
if !found {
|
||||
t.Fatalf("part should be found after replace")
|
||||
}
|
||||
|
||||
if foundPart.MD5 != "md5_2" {
|
||||
t.Fatalf("part should be replaced with new MD5")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_UpdateOrStore tests UpdateOrStore method
|
||||
func TestMUParts_UpdateOrStore(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
part1 := &Part{
|
||||
ID: 1,
|
||||
MD5: "md5_1",
|
||||
Size: 1024,
|
||||
Inode: 100,
|
||||
UploadTime: time.Now(),
|
||||
}
|
||||
|
||||
part2 := &Part{
|
||||
ID: 1,
|
||||
MD5: "md5_2",
|
||||
Size: 2048,
|
||||
Inode: 200,
|
||||
UploadTime: time.Now().Add(time.Hour), // Later time
|
||||
}
|
||||
|
||||
part3 := &Part{
|
||||
ID: 1,
|
||||
MD5: "md5_3",
|
||||
Size: 4096,
|
||||
Inode: 200, // Same inode as part1
|
||||
UploadTime: time.Now().Add(2 * time.Hour),
|
||||
}
|
||||
|
||||
// First insert
|
||||
_, updated, conflict := parts.UpdateOrStore(part1)
|
||||
if updated || conflict {
|
||||
t.Fatalf("first insert should not be update or conflict")
|
||||
}
|
||||
|
||||
// Update with different inode and later time
|
||||
oldInode, updated, conflict := parts.UpdateOrStore(part2)
|
||||
if !updated || conflict {
|
||||
t.Fatalf("update with later time should succeed")
|
||||
}
|
||||
if oldInode != 100 {
|
||||
t.Fatalf("old inode should be 100, got %d", oldInode)
|
||||
}
|
||||
|
||||
// Try to update with same inode - this should not update and not conflict
|
||||
oldInode, updated, conflict = parts.UpdateOrStore(part3)
|
||||
if updated || conflict {
|
||||
t.Fatalf("update with same inode should not update or conflict, got updated=%v conflict=%v", updated, conflict)
|
||||
}
|
||||
if oldInode != 200 {
|
||||
t.Fatalf("old inode should be 200 (from part2), got %d", oldInode)
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_Hash tests Hash method
|
||||
func TestMUParts_Hash(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
part1 := &Part{ID: 1, MD5: "md5_1", Size: 1024, Inode: 100}
|
||||
part2 := &Part{ID: 2, MD5: "md5_2", Size: 2048, Inode: 200}
|
||||
|
||||
parts.Insert(part1, false)
|
||||
parts.Insert(part2, false)
|
||||
|
||||
// Test hash for existing part
|
||||
if !parts.Hash(part1) {
|
||||
t.Fatalf("hash should find existing part")
|
||||
}
|
||||
|
||||
// Test hash for non-existing part
|
||||
part3 := &Part{ID: 3, MD5: "md5_3", Size: 4096, Inode: 300}
|
||||
if parts.Hash(part3) {
|
||||
t.Fatalf("hash should not find non-existing part")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUMultipartExtend_Empty tests empty MultipartExtend
|
||||
func TestMUMultipartExtend_Empty(t *testing.T) {
|
||||
me := NewMultipartExtend()
|
||||
|
||||
bytes, err := me.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("empty extend bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
me2 := MultipartExtendFromBytes(bytes)
|
||||
// According to the implementation, empty extend returns nil
|
||||
if me2 != nil {
|
||||
t.Fatalf("deserialized empty extend should be nil, got %v", me2)
|
||||
}
|
||||
|
||||
// Test with non-empty extend to ensure the method works
|
||||
me["test"] = "value"
|
||||
bytes, err = me.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("non-empty extend bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
me3 := MultipartExtendFromBytes(bytes)
|
||||
if me3 == nil {
|
||||
t.Fatalf("deserialized non-empty extend should not be nil")
|
||||
}
|
||||
|
||||
if len(me3) != 1 || me3["test"] != "value" {
|
||||
t.Fatalf("deserialized extend should have correct content")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUMultipartExtend_Large tests MultipartExtend with large data
|
||||
func TestMUMultipartExtend_Large(t *testing.T) {
|
||||
me := NewMultipartExtend()
|
||||
|
||||
// Add many key-value pairs
|
||||
for i := 0; i < 1000; i++ {
|
||||
key := util.RandomString(20, util.UpperLetter|util.LowerLetter|util.Numeric)
|
||||
value := util.RandomString(100, util.UpperLetter|util.LowerLetter|util.Numeric)
|
||||
me[key] = value
|
||||
}
|
||||
|
||||
bytes, err := me.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("large extend bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
me2 := MultipartExtendFromBytes(bytes)
|
||||
if !reflect.DeepEqual(me, me2) {
|
||||
t.Fatalf("large extend deserialization mismatch")
|
||||
}
|
||||
|
||||
t.Logf("large extend encoded length: %v", len(bytes))
|
||||
}
|
||||
|
||||
// TestMUMultipart_Concurrent tests concurrent access to Multipart
|
||||
func TestMUMultipart_Concurrent(t *testing.T) {
|
||||
multipart := &Multipart{
|
||||
id: "test-session",
|
||||
key: "test-key",
|
||||
initTime: time.Now(),
|
||||
parts: PartsFromBytes(nil),
|
||||
}
|
||||
|
||||
// Concurrent goroutines
|
||||
done := make(chan bool, 10)
|
||||
|
||||
for i := 0; i < 10; i++ {
|
||||
go func(id int) {
|
||||
part := &Part{
|
||||
ID: uint16(id),
|
||||
MD5: util.RandomString(16, util.UpperLetter|util.Numeric),
|
||||
Size: uint64(1024 + id),
|
||||
Inode: uint64(1000 + id),
|
||||
UploadTime: time.Now(),
|
||||
}
|
||||
|
||||
multipart.InsertPart(part, false)
|
||||
multipart.Parts() // Test read access
|
||||
done <- true
|
||||
}(i)
|
||||
}
|
||||
|
||||
// Wait for all goroutines
|
||||
for i := 0; i < 10; i++ {
|
||||
<-done
|
||||
}
|
||||
|
||||
parts := multipart.Parts()
|
||||
if len(parts) != 10 {
|
||||
t.Fatalf("expected 10 parts, got %d", len(parts))
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUMultipart_EdgeCases tests edge cases for Multipart
|
||||
func TestMUMultipart_EdgeCases(t *testing.T) {
|
||||
// Test with nil parts
|
||||
multipart := &Multipart{
|
||||
id: "test-session",
|
||||
key: "test-key",
|
||||
initTime: time.Now(),
|
||||
parts: nil,
|
||||
}
|
||||
|
||||
// This should initialize parts
|
||||
multipart.InsertPart(&Part{ID: 1, MD5: "test", Size: 1024, Inode: 100}, false)
|
||||
|
||||
parts := multipart.Parts()
|
||||
if len(parts) != 1 {
|
||||
t.Fatalf("nil parts should be initialized, expected 1 part, got %d", len(parts))
|
||||
}
|
||||
|
||||
// Test UpdateOrStorePart with nil parts
|
||||
multipart2 := &Multipart{
|
||||
id: "test-session-2",
|
||||
key: "test-key-2",
|
||||
initTime: time.Now(),
|
||||
parts: nil,
|
||||
}
|
||||
|
||||
oldInode, updated, conflict := multipart2.UpdateOrStorePart(&Part{ID: 1, MD5: "test", Size: 1024, Inode: 100})
|
||||
if updated || conflict {
|
||||
t.Fatalf("first UpdateOrStorePart should not be update or conflict")
|
||||
}
|
||||
if oldInode != 0 {
|
||||
t.Fatalf("old inode should be 0 for new part, got %d", oldInode)
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUMultipart_BytesEmpty tests Multipart serialization with empty data
|
||||
func TestMUMultipart_BytesEmpty(t *testing.T) {
|
||||
multipart := &Multipart{
|
||||
id: "",
|
||||
key: "",
|
||||
initTime: time.Time{},
|
||||
parts: PartsFromBytes(nil),
|
||||
extend: NewMultipartExtend(),
|
||||
}
|
||||
|
||||
bytes, err := multipart.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("empty multipart bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
multipart2 := MultipartFromBytes(bytes)
|
||||
if multipart2.id != "" || multipart2.key != "" {
|
||||
t.Fatalf("deserialized empty multipart should have empty id and key")
|
||||
}
|
||||
|
||||
if multipart2.parts.Len() != 0 {
|
||||
t.Fatalf("deserialized empty multipart should have empty parts")
|
||||
}
|
||||
|
||||
if len(multipart2.extend) != 0 {
|
||||
t.Fatalf("deserialized empty multipart should have empty extend")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_Sort tests sorting functionality
|
||||
func TestMUParts_Sort(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
// Insert parts in reverse order
|
||||
for i := 10; i >= 1; i-- {
|
||||
part := &Part{
|
||||
ID: uint16(i),
|
||||
MD5: util.RandomString(16, util.UpperLetter|util.Numeric),
|
||||
Size: uint64(1024 * i),
|
||||
Inode: uint64(1000 + i),
|
||||
}
|
||||
parts.Insert(part, false)
|
||||
}
|
||||
|
||||
// Verify parts are sorted by ID
|
||||
for i := 0; i < parts.Len()-1; i++ {
|
||||
if parts[i].ID >= parts[i+1].ID {
|
||||
t.Fatalf("parts should be sorted by ID, found %d >= %d", parts[i].ID, parts[i+1].ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUParts_RemoveNonExistent tests removing non-existent parts
|
||||
func TestMUParts_RemoveNonExistent(t *testing.T) {
|
||||
parts := PartsFromBytes(nil)
|
||||
|
||||
// Add some parts
|
||||
for i := 1; i <= 5; i++ {
|
||||
part := &Part{
|
||||
ID: uint16(i),
|
||||
MD5: util.RandomString(16, util.UpperLetter|util.Numeric),
|
||||
Size: uint64(1024 * i),
|
||||
Inode: uint64(1000 + i),
|
||||
}
|
||||
parts.Insert(part, false)
|
||||
}
|
||||
|
||||
originalLen := parts.Len()
|
||||
|
||||
// Try to remove non-existent part
|
||||
parts.Remove(999)
|
||||
|
||||
if parts.Len() != originalLen {
|
||||
t.Fatalf("removing non-existent part should not change length")
|
||||
}
|
||||
}
|
||||
|
||||
// TestMUMultipartExtend_SpecialCharacters tests MultipartExtend with special characters
|
||||
func TestMUMultipartExtend_SpecialCharacters(t *testing.T) {
|
||||
me := NewMultipartExtend()
|
||||
|
||||
// Add keys and values with special characters
|
||||
me["key with spaces"] = "value with spaces"
|
||||
me["key:with:colons"] = "value:with:colons"
|
||||
me["key=with=equals"] = "value=with=equals"
|
||||
me["key&with&ersands"] = "value&with&ersands"
|
||||
me["key\nwith\nnewlines"] = "value\nwith\nnewlines"
|
||||
me["key\twith\ttabs"] = "value\twith\ttabs"
|
||||
|
||||
bytes, err := me.Bytes()
|
||||
if err != nil {
|
||||
t.Fatalf("special characters extend bytes should not fail: %v", err)
|
||||
}
|
||||
|
||||
me2 := MultipartExtendFromBytes(bytes)
|
||||
if !reflect.DeepEqual(me, me2) {
|
||||
t.Fatalf("special characters extend deserialization mismatch")
|
||||
}
|
||||
}
|
||||
|
||||
@ -483,7 +483,7 @@ func TestSortedExtentEdgeCases(t *testing.T) {
|
||||
assert.NotPanics(t, func() {
|
||||
se.IsEmpty()
|
||||
se.Len()
|
||||
se.String()
|
||||
_ = se.String()
|
||||
se.LayerSize()
|
||||
se.Size()
|
||||
})
|
||||
|
||||
@ -305,7 +305,7 @@ func TestSortedObjExtentsConcurrentAccess(t *testing.T) {
|
||||
se.IsEmpty()
|
||||
se.Size()
|
||||
se.LayerSize()
|
||||
se.String()
|
||||
_ = se.String()
|
||||
done <- true
|
||||
}()
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user