test(data): add testcase for datanode op log when disk is full.

Signed-off-by: Victor1319 <zengxuewei@oppo.com>
This commit is contained in:
Victor1319 2024-09-23 15:47:58 +08:00 committed by chihe
parent a5214f0478
commit 7f43f724ec
3 changed files with 191 additions and 49 deletions

View File

@ -15,12 +15,20 @@ import (
func (d *DataNode) startStat(cfg *config.Config) {
logDir := cfg.GetString(ConfigKeyLogDir)
var err error
stat.DpStat, err = stat.NewOpLogger(logDir, "dp_op.log", stat.DefaultMaxOps, stat.DefaultDuration)
logLeftSpaceLimitRatioStr := cfg.GetString("logLeftSpaceLimitRatio")
logLeftSpaceLimitRatio, err := strconv.ParseFloat(logLeftSpaceLimitRatioStr, 64)
if err != nil {
log.LogErrorf("get log limt ratio failed, err %s", err.Error())
logLeftSpaceLimitRatio = log.DefaultLogLeftSpaceLimitRatio
}
stat.DpStat, err = stat.NewOpLogger(logDir, "dp_op.log", stat.DefaultMaxOps, stat.DefaultDuration, logLeftSpaceLimitRatio)
if err != nil {
log.LogErrorf("DpStat init failed.err:%+v", err)
return
}
stat.DiskStat, err = stat.NewOpLogger(logDir, "disk_op.log", stat.DefaultMaxOps, stat.DefaultDuration)
stat.DiskStat, err = stat.NewOpLogger(logDir, "disk_op.log", stat.DefaultMaxOps, stat.DefaultDuration, logLeftSpaceLimitRatio)
if err != nil {
log.LogErrorf("DiskStat init failed.err:%+v", err)
return

View File

@ -26,9 +26,17 @@ import (
"sync/atomic"
"time"
"github.com/cubefs/cubefs/util/fileutil"
"github.com/cubefs/cubefs/util/log"
)
const (
DefaultMaxOps = 100
DefaultDuration = time.Minute
defaultSep = "+"
oplogModule = "oplogs"
)
var (
DpStat = new(OpLogger)
DiskStat = new(OpLogger)
@ -55,16 +63,10 @@ type OpLogger struct {
reserveTime time.Duration
dir string
filename string
leftSpace int64
}
const (
DefaultMaxOps = 100
DefaultDuration = time.Minute
defaultSep = "+"
oplogModule = "oplogs"
)
func NewOpLogger(dir, filename string, maxOps int, duration time.Duration) (*OpLogger, error) {
func NewOpLogger(dir, filename string, maxOps int, duration time.Duration, leftSpaceRatio float64) (*OpLogger, error) {
dir = path.Join(dir, oplogModule)
fi, err := os.Stat(dir)
if err != nil {
@ -74,6 +76,18 @@ func NewOpLogger(dir, filename string, maxOps int, duration time.Duration) (*OpL
return new(OpLogger), errors.New(dir + " is not a directory")
}
}
fs, err := fileutil.Statfs(dir)
if err != nil {
return nil, errors.New("statfs dir failed, " + err.Error())
}
if leftSpaceRatio < log.DefaultLogLeftSpaceLimitRatio {
leftSpaceRatio = log.DefaultLogLeftSpaceLimitRatio
}
logLeftSpace := int64(float64((fs.Blocks * uint64(fs.Bsize))) * leftSpaceRatio)
_ = os.Chmod(dir, 0o755)
logger := &OpLogger{
opCounts: map[string]*int32{},
@ -88,7 +102,9 @@ func NewOpLogger(dir, filename string, maxOps int, duration time.Duration) (*OpL
reserveTime: MaxReservedDays,
dir: dir,
filename: filename,
leftSpace: logLeftSpace,
}
go logger.startFlushing()
return logger, nil
}
@ -102,26 +118,14 @@ func RecordStat(partitionID uint64, op string, dataPath string) {
}
}
func (l *OpLogger) Record(name string) {
l.RecordOp(name, "")
// used for test
func (l *OpLogger) SetArgs(reserveTime time.Duration, fileSize int64) {
l.reserveTime = reserveTime
l.fileSize = fileSize
}
func (l *OpLogger) incrementCount(counts map[string]*int32, key string) {
l.RLock()
if _, ok := counts[key]; !ok {
l.RUnlock()
l.Lock()
if _, ok = counts[key]; !ok {
counts[key] = new(int32)
}
atomic.AddInt32(counts[key], 1)
l.Unlock()
return
}
atomic.AddInt32(counts[key], 1)
l.RUnlock()
func (l *OpLogger) Record(name string) {
l.RecordOp(name, "")
}
func (l *OpLogger) RecordOp(name, op string) {
@ -158,10 +162,6 @@ func (l *OpLogger) SetFileSize(fileSize int64) {
l.fileSize = fileSize
}
func (l *OpLogger) SetReserveTime(duration time.Duration) {
l.reserveTime = duration
}
func (l *OpLogger) GetMasterOps() []*Operation {
l.Lock()
defer l.Unlock()
@ -174,6 +174,30 @@ func (l *OpLogger) GetPrevOps() []*Operation {
return l.getAllOps(l.opCountsPrev)
}
func (l *OpLogger) Close() {
l.ticker.Stop()
l.done <- true
l.flush()
}
func (l *OpLogger) incrementCount(counts map[string]*int32, key string) {
l.RLock()
if _, ok := counts[key]; !ok {
l.RUnlock()
l.Lock()
if _, ok = counts[key]; !ok {
counts[key] = new(int32)
}
atomic.AddInt32(counts[key], 1)
l.Unlock()
return
}
atomic.AddInt32(counts[key], 1)
l.RUnlock()
}
func (l *OpLogger) startFlushing() {
for {
select {
@ -233,23 +257,48 @@ func (l *OpLogger) remove() {
if err != nil {
return
}
for _, entry := range entries {
fileInfo, _ := entry.Info()
if fileInfo == nil {
oldLogs := make([]os.DirEntry, 0)
for _, e := range entries {
if strings.HasPrefix(e.Name(), l.filename) && strings.HasSuffix(e.Name(), ShiftedExtension) {
oldLogs = append(oldLogs, e)
}
}
if len(oldLogs) == 0 {
return
}
sort.Slice(oldLogs, func(i, j int) bool {
return oldLogs[i].Name() < oldLogs[j].Name()
})
fs, err := fileutil.Statfs(l.dir)
if err != nil {
log.LogErrorf("remove stat fs failed, err %v", err)
return
}
leftSpace := int64(fs.Bavail*uint64(fs.Bsize)) - l.leftSpace
for _, e := range oldLogs {
info, err := e.Info()
if err != nil {
log.LogErrorf("get log info failed, file %s, err %s", e.Name(), err.Error())
continue
}
if fileInfo.IsDir() {
if time.Since(info.ModTime()) < l.reserveTime && leftSpace > 0 {
continue
}
if !strings.HasPrefix(fileInfo.Name(), l.filename) {
leftSpace += info.Size()
delFile := path.Join(l.dir, e.Name())
err = os.Remove(delFile)
if err != nil && !os.IsNotExist(err) {
log.LogErrorf("delete file failed, path %s, err %s", delFile, err.Error())
continue
}
if !strings.HasSuffix(fileInfo.Name(), ShiftedExtension) {
continue
}
if time.Since(fileInfo.ModTime()) > l.reserveTime {
os.Remove(path.Join(l.dir, fileInfo.Name()))
}
}
}
@ -279,9 +328,3 @@ func (l *OpLogger) getAllOps(m map[string]*int32) []*Operation {
})
return ops
}
func (l *OpLogger) Close() {
l.ticker.Stop()
l.done <- true
l.flush()
}

View File

@ -15,12 +15,14 @@
package stat
import (
"fmt"
"os"
"path"
"testing"
"time"
"github.com/cubefs/cubefs/util/errors"
"github.com/cubefs/cubefs/util/fileutil"
"github.com/stretchr/testify/require"
)
@ -45,3 +47,92 @@ func TestStatistic(t *testing.T) {
EndStat("test2", err, bgTime, 100)
time.Sleep(3 * time.Second)
}
func TestDataNodeOpStatsRotate(t *testing.T) {
tmpDir, err := os.MkdirTemp("", "test_op_*")
require.NoError(t, err)
defer func() {
os.RemoveAll(tmpDir)
}()
logName := "test_op.log"
fileSize := 1024 * 1024
count := 30240 // every file is 1.6M
tickCnt := 10 // generate 10 files
expireInterval := 6
opLog, err := NewOpLogger(tmpDir, logName, 0, time.Second, 0)
require.NoError(t, err)
// set rotate size as 1M
opLog.SetArgs(time.Second*time.Duration(expireInterval), int64(fileSize))
start := time.Now()
for {
if time.Since(start) > time.Second*time.Duration(tickCnt) {
break
}
for i := 0; i < count; i++ {
key := fmt.Sprintf("op_name_%d", i)
msg := fmt.Sprintf("tt_op_xxxx_xxx_%d", i)
opLog.RecordOp(key, msg)
}
}
items, err := os.ReadDir(opLog.dir)
require.NoError(t, err)
for _, e := range items {
t.Logf("got log files, e %s", e.Name())
}
// assert there are rotate files, and expire is vaild.
require.True(t, len(items) > 1 && len(items) <= expireInterval+1)
}
func TestDataNodeOpStatsDiskFull(t *testing.T) {
tmpDir, err := os.MkdirTemp("", "test_op_*")
require.NoError(t, err)
defer func() {
os.RemoveAll(tmpDir)
}()
logName := "test_op.log"
fileSize := 1024 * 1024 // rotate over 1m
count := 30240 // every file is 1.6M
leftCnt := 4
tickCnt := 10
opLog, _ := NewOpLogger(tmpDir, logName, 0, time.Second, 0)
require.NoError(t, err)
fs, err := fileutil.Statfs(tmpDir)
require.NoError(t, err)
avail := fs.Bavail * uint64(fs.Bsize)
opLog.leftSpace = int64(avail) + int64(leftCnt*fileSize)
// set rotate size as 1M
opLog.SetArgs(time.Hour, int64(fileSize))
start := time.Now()
for {
if time.Since(start) > time.Second*time.Duration(tickCnt) {
break
}
for i := 0; i < count; i++ {
key := fmt.Sprintf("op_name_%d", i)
msg := fmt.Sprintf("tt_op_xxxx_xxx_%d", i)
opLog.RecordOp(key, msg)
}
}
items, err := os.ReadDir(opLog.dir)
require.NoError(t, err)
for _, e := range items {
t.Logf("got log files, e %s", e.Name())
}
require.True(t, len(items) <= 3)
}