From c2611b2459b003e0b2bc468abb31102169ba48d3 Mon Sep 17 00:00:00 2001 From: zhumingze Date: Wed, 7 May 2025 17:16:37 +0800 Subject: [PATCH] feat(data): Check dp availability at random times. #1000022788 Signed-off-by: zhumingze --- datanode/partition.go | 21 ++++++- datanode/partition_test.go | 110 +++++++++++++++++++++++++++++++++++++ 2 files changed, 129 insertions(+), 2 deletions(-) create mode 100644 datanode/partition_test.go diff --git a/datanode/partition.go b/datanode/partition.go index 1cdd4f25f..91e60284b 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -20,6 +20,7 @@ import ( "hash/crc32" "io/ioutil" "math" + "math/rand" "net" "os" "path" @@ -57,6 +58,13 @@ const ( RaftStatusRunning = 1 ) +const ( + DpCheckBaseInterval = 7200 + DpCheckRandomRange = 1800 + DpMinCheckInterval = DpCheckBaseInterval - DpCheckRandomRange + DpMaxCheckInterval = DpCheckBaseInterval + DpCheckRandomRange +) + type DataPartitionMetadata struct { VolumeID string PartitionID uint64 @@ -866,7 +874,13 @@ func (dp *DataPartition) statusUpdateScheduler() { ticker := time.NewTicker(time.Minute) snapshotTicker := time.NewTicker(time.Minute * 5) peersTicker := time.NewTicker(10 * time.Second) - dpCheckTicket := time.NewTicker(2 * time.Hour) + + genCheckInterval := func() time.Duration { + interval := DpMinCheckInterval + rand.Intn(DpMaxCheckInterval-DpMinCheckInterval+1) + return time.Duration(interval) * time.Second + } + dpCheckTimer := time.NewTimer(genCheckInterval()) + var index int for { select { @@ -887,11 +901,14 @@ func (dp *DataPartition) statusUpdateScheduler() { dp.ReloadSnapshot() case <-peersTicker.C: dp.validatePeers() - case <-dpCheckTicket.C: + case <-dpCheckTimer.C: dp.checkAvailable() + dpCheckTimer.Reset(genCheckInterval()) case <-dp.stopC: ticker.Stop() snapshotTicker.Stop() + peersTicker.Stop() + dpCheckTimer.Stop() return } } diff --git a/datanode/partition_test.go b/datanode/partition_test.go new file mode 100644 index 000000000..d8d8d51e6 --- /dev/null +++ b/datanode/partition_test.go @@ -0,0 +1,110 @@ +// Copyright 2025 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + +package datanode + +import ( + "os" + "path" + "path/filepath" + "syscall" + "testing" + + "github.com/cubefs/cubefs/datanode/storage" + "github.com/cubefs/cubefs/proto" + "github.com/cubefs/cubefs/util" + "github.com/stretchr/testify/require" +) + +func TestCheckAvailable(t *testing.T) { + tmpDir, err := os.MkdirTemp(".", "") + defer os.RemoveAll(tmpDir) + require.NoError(t, err) + + testDpPath := path.Join(tmpDir, "dp1") + err = os.Mkdir(testDpPath, 0o755) + require.NoError(t, err) + + dp := &DataPartition{ + disk: &Disk{dataNode: &DataNode{}}, + partitionID: 1, + partitionSize: 1 * util.TB, + path: testDpPath, + partitionStatus: proto.ReadWrite, + diskErrCnt: 0, + config: &dataPartitionCfg{}, + volVersionInfoList: &proto.VolVersionInfoList{}, + extentStore: &storage.ExtentStore{}, + } + + statusFilePath := filepath.Join(dp.path, DpStatusFile) + + t.Run("file not exist", func(t *testing.T) { + err := dp.checkAvailable() + require.Error(t, err) + require.Equal(t, proto.ReadWrite, dp.partitionStatus) + require.Equal(t, uint64(0), dp.diskErrCnt) + }) + + require.NoError(t, os.WriteFile(statusFilePath, []byte(DpStatusFile), 0o644)) + + t.Run("normal readWrite", func(t *testing.T) { + err := dp.checkAvailable() + require.NoError(t, err) + require.Equal(t, proto.ReadWrite, dp.partitionStatus) + require.Equal(t, uint64(0), dp.diskErrCnt) + }) + + t.Run("disk err simulate", func(t *testing.T) { + fp, err := os.OpenFile(statusFilePath, os.O_TRUNC|os.O_RDWR, 0o755) + require.NoError(t, err) + defer fp.Close() + + data := []byte(DpStatusFile) + t.Run("read err", func(t *testing.T) { + _, err = fp.WriteAt(data, 0) + require.NoError(t, err) + fp.Close() + _, err = fp.ReadAt(data, 0) + require.Error(t, err) + err = syscall.EIO + dp.checkIsDiskError(err, ReadFlag) + require.Equal(t, proto.Unavailable, dp.partitionStatus) + require.Equal(t, uint64(1), dp.diskErrCnt) + }) + + fp, err = os.OpenFile(statusFilePath, os.O_TRUNC|os.O_RDWR, 0o755) + + t.Run("sync err", func(t *testing.T) { + _, err = fp.WriteAt(data, 0) + require.NoError(t, err) + fp.Close() + err = fp.Sync() + require.Error(t, err) + err = syscall.EIO + dp.checkIsDiskError(err, ReadFlag) + require.Equal(t, proto.Unavailable, dp.partitionStatus) + require.Equal(t, uint64(2), dp.diskErrCnt) + }) + + t.Run("write err", func(t *testing.T) { + _, err = fp.WriteAt(data, 0) + require.Error(t, err) + err = syscall.EIO + dp.checkIsDiskError(err, WriteFlag) + require.Equal(t, proto.Unavailable, dp.partitionStatus) + require.Equal(t, uint64(3), dp.diskErrCnt) + }) + }) +}