diff --git a/datanode/stat_log.go b/datanode/stat_log.go index 3dceb9977..54d9216fd 100644 --- a/datanode/stat_log.go +++ b/datanode/stat_log.go @@ -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 diff --git a/util/stat/datanode_op_stats.go b/util/stat/datanode_op_stats.go index 14acc8526..9f9dff205 100644 --- a/util/stat/datanode_op_stats.go +++ b/util/stat/datanode_op_stats.go @@ -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() -} diff --git a/util/stat/statistic_test.go b/util/stat/statistic_test.go index 4eda83fb0..ea9f67810 100644 --- a/util/stat/statistic_test.go +++ b/util/stat/statistic_test.go @@ -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) +}