feat(datamaster): Supplement test cases related to lost and bad disks. #1000004886

Signed-off-by: zhumingze <zhumingze@oppo.com>
This commit is contained in:
zhumingze 2025-04-02 10:12:55 +08:00 committed by zhumingze1108
parent 1964c60240
commit bfe93e050c
8 changed files with 211 additions and 3 deletions

View File

@ -639,6 +639,9 @@ type dpLoadInfo struct {
// RestorePartition reads the files stored on the local disk and restores the data partitions.
func (d *Disk) RestorePartition(visitor PartitionVisitor) (err error) {
if d.space.dataNode.localServerAddr == "" {
return
}
convert := func(node *proto.DataNodeInfo) *DataNodeInfo {
result := &DataNodeInfo{}
result.Addr = node.Addr

View File

@ -344,7 +344,9 @@ func (manager *SpaceManager) reloadDisk(path string) (err error) {
manager.putDisk(disk)
return
}
log.LogWarnf("[reloadDisk] Successfully reloaded: %s", params.path)
msg := fmt.Sprintf("Successfully reloaded: %s", params.path)
auditlog.LogDataNodeOp("ReloadDisk", msg, nil)
log.LogWarnf("[reloadDisk] %s", msg)
}(diskParams)
return
}
@ -503,7 +505,7 @@ func (manager *SpaceManager) GetDisk(path string) (d *Disk, err error) {
d = disk
return
}
err = fmt.Errorf("disk(%v) not exsit", path)
err = fmt.Errorf("disk(%v) not exist", path)
return
}

View File

@ -2276,7 +2276,7 @@ func (s *DataNode) handlePacketToReloadDisk(p *repl.Packet) {
if err != nil {
return
}
log.LogDebugf("action[handlePacketToReloadDisk] try reload disk %v req %v", request.DiskPath, task.RequestID)
log.LogWarnf("action[handlePacketToReloadDisk] try reload disk %v req %v", request.DiskPath, task.RequestID)
err = s.space.reloadDisk(request.DiskPath)
if err != nil {

View File

@ -15,13 +15,18 @@
package datanode
import (
"encoding/json"
"os"
"path"
"sync"
"testing"
"time"
"github.com/cubefs/cubefs/datanode/repl"
"github.com/cubefs/cubefs/datanode/storage"
"github.com/cubefs/cubefs/proto"
"github.com/cubefs/cubefs/util"
"github.com/cubefs/cubefs/util/atomicutil"
"github.com/stretchr/testify/require"
"golang.org/x/time/rate"
)
@ -110,3 +115,150 @@ func TestSkipAppendWrite(t *testing.T) {
t.Logf("handle write packet, result code(%v)", p.ResultCode)
require.EqualValues(t, proto.OpArgMismatchErr, p.ResultCode)
}
func newPacketForTest(task *proto.AdminTask) *repl.Packet {
data, _ := json.Marshal(task)
return &repl.Packet{
Packet: proto.Packet{
Data: data,
},
}
}
func TestDeleteLostDisk(t *testing.T) {
dn := &DataNode{
space: &SpaceManager{
disks: make(map[string]*Disk),
diskList: []string{},
diskUtils: make(map[string]*atomicutil.Float64),
},
}
testDiskPath := "/test/disk1"
lostDisk := NewLostDisk(
testDiskPath,
1*util.TB,
0,
3,
dn.space,
true,
)
dn.space.putDisk(lostDisk)
t.Run("normal delete disk", func(t *testing.T) {
req := &proto.DeleteLostDiskRequest{DiskPath: testDiskPath}
task := &proto.AdminTask{
OpCode: proto.OpDeleteLostDisk,
Request: req,
}
p := newPacketForTest(task)
dn.handlePacketToDeleteLostDisk(p)
require.Equal(t, proto.OpOk, p.ResultCode)
_, err := dn.space.GetDisk(testDiskPath)
require.Error(t, err)
require.Contains(t, err.Error(), "not exist")
})
t.Run("delete unexist disk", func(t *testing.T) {
invalidReq := &proto.DeleteLostDiskRequest{DiskPath: "/invalid/path"}
task := &proto.AdminTask{
OpCode: proto.OpDeleteLostDisk,
Request: invalidReq,
}
p := newPacketForTest(task)
dn.handlePacketToDeleteLostDisk(p)
require.Equal(t, proto.OpIntraGroupNetErr, p.ResultCode)
require.Contains(t, string(p.Data), "not exist")
})
}
func TestReloadDisk(t *testing.T) {
tmpDir, err := os.MkdirTemp(".", "")
defer os.RemoveAll(tmpDir)
require.NoError(t, err)
dn := &DataNode{
diskReadFlow: 1 * util.MB,
diskWriteFlow: 1 * util.MB,
diskReadIocc: 10,
diskWriteIocc: 10,
}
sm := &SpaceManager{
disks: make(map[string]*Disk),
dataNode: dn,
diskList: []string{},
diskUtils: make(map[string]*atomicutil.Float64),
}
dn.space = sm
testDiskPath := path.Join(tmpDir, "disk1")
err = os.Mkdir(testDiskPath, 0o755)
require.NoError(t, err)
disk := NewLostDisk(
testDiskPath,
1*util.TB,
0,
3,
dn.space,
true,
)
dn.space.putDisk(disk)
req := &proto.ReloadDiskRequest{DiskPath: testDiskPath}
task := &proto.AdminTask{
OpCode: proto.OpReloadDisk,
Request: req,
}
t.Run("normal reload disk", func(t *testing.T) {
p := newPacketForTest(task)
dn.handlePacketToReloadDisk(p)
require.Equal(t, proto.OpOk, p.ResultCode)
require.Eventually(t, func() bool {
disk, _ := dn.space.GetDisk(testDiskPath)
return disk != nil && !disk.isLost
}, 3*time.Second, 100*time.Millisecond, "disk not loaded")
})
t.Run("reload conflict", func(t *testing.T) {
var wg sync.WaitGroup
results := make(chan uint8, 2)
for i := 0; i < 2; i++ {
wg.Add(1)
go func() {
defer wg.Done()
p := newPacketForTest(&proto.AdminTask{
OpCode: proto.OpReloadDisk,
Request: &proto.ReloadDiskRequest{DiskPath: testDiskPath},
})
dn.handlePacketToReloadDisk(p)
results <- p.ResultCode
}()
}
go func() {
wg.Wait()
close(results)
}()
var successCount, errorCount int
for res := range results {
switch res {
case proto.OpOk:
successCount++
case proto.OpIntraGroupNetErr:
errorCount++
}
}
require.Equal(t, 1, successCount)
require.Equal(t, 1, errorCount)
})
}

View File

@ -1876,3 +1876,19 @@ func TestCleanEmptyMetaPartition(t *testing.T) {
process(reqUrl, t)
}
func TestDeleteLostDisk(t *testing.T) {
addr := mds5Addr
disk := "/cfs/disk2"
reqUrl := fmt.Sprintf("%v%v?addr=%v&disk=%v", hostAddr, proto.DeleteLostDisk, addr, disk)
process(reqUrl, t)
}
func TestReloadDisk(t *testing.T) {
addr := mds5Addr
disk := "/cfs/disk"
reqUrl := fmt.Sprintf("%v%v?addr=%v&disk=%v", hostAddr, proto.ReloadDisk, addr, disk)
process(reqUrl, t)
}

View File

@ -6,6 +6,7 @@ import (
"time"
"github.com/cubefs/cubefs/proto"
"github.com/stretchr/testify/require"
)
func TestDataNode(t *testing.T) {
@ -20,6 +21,7 @@ func TestDataNode(t *testing.T) {
server.cluster.checkDataNodeHeartbeat()
time.Sleep(5 * time.Second)
getDataNodeInfo(addr, t)
updateDisks(addr, t)
decommissionDataNode(addr, t)
for i := 0; i < 10; i++ { // decommission is async process
_, err = server.cluster.dataNode(addr)
@ -44,3 +46,16 @@ func decommissionDataNode(addr string, t *testing.T) {
reqURL := fmt.Sprintf("%v%v?addr=%v", hostAddr, proto.DecommissionDataNode, addr)
process(reqURL, t)
}
func updateDisks(addr string, t *testing.T) {
dn, err := server.cluster.dataNode(addr)
require.NoError(t, err)
dn.AllDisks = []string{"/data1"}
allDisk := []string{"/data1", "/data2", "/data3"}
badDisk := []string{"/data1"}
updated, _ := dn.updateDisks(allDisk, badDisk)
require.Equal(t, updated, true)
require.Equal(t, allDisk, dn.AllDisks)
require.Equal(t, badDisk, dn.BadDisks)
}

View File

@ -147,6 +147,12 @@ func (mds *MockDataServer) serveConn(rc net.Conn) {
case proto.OpDataPartitionTryToLeader:
err = mds.handleTryToLeader(conn, req, adminTask)
Printf("data node [%v] try to leader,id[%v],err:%v\n", mds.TcpAddr, adminTask.ID, err)
case proto.OpDeleteLostDisk:
err = mds.handleDeleteLostDisk(conn, req, adminTask)
Printf("data node [%v] try to delete disk[%v],err:%v\n", mds.TcpAddr, adminTask.Disk, err)
case proto.OpReloadDisk:
err = mds.handleReloadDisk(conn, req, adminTask)
Printf("data node [%v] try to delete disk[%v],err:%v\n", mds.TcpAddr, adminTask.Disk, err)
default:
fmt.Printf("unknown code [%v]\n", req.Opcode)
}
@ -297,6 +303,9 @@ func (mds *MockDataServer) handleHeartbeats(conn net.Conn, pkg *proto.Packet, ta
response.AllDisks = []string{
"/cfs/disk",
}
response.LostDisks = []string{
"/cfs/disk2",
}
mds.RLock()
for _, partition := range mds.partitions {
@ -396,3 +405,13 @@ func buildSnapshot() (files []*proto.File) {
files = append(files, f3)
return
}
func (mds *MockDataServer) handleDeleteLostDisk(conn net.Conn, pkg *proto.Packet, adminTask *proto.AdminTask) (err error) {
err = responseAckOKToMaster(conn, pkg, nil)
return
}
func (mds *MockDataServer) handleReloadDisk(conn net.Conn, pkg *proto.Packet, adminTask *proto.AdminTask) (err error) {
err = responseAckOKToMaster(conn, pkg, nil)
return
}

View File

@ -36,6 +36,7 @@ const (
type AdminTask struct {
ID string
PartitionID uint64
Disk string
OpCode uint8
OperatorAddr string
Status int8