fix(master): do not delete decommission disk from list when recover task is submitted

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-07-09 22:41:56 +08:00 committed by AmazingChi
parent d479a79f73
commit a8d8c96072
8 changed files with 114 additions and 68 deletions

View File

@ -22,6 +22,7 @@ import (
"net"
"os"
"path"
"runtime/debug"
"strconv"
"strings"
"sync/atomic"

View File

@ -37,7 +37,6 @@ import (
"github.com/cubefs/cubefs/util/errors"
"github.com/cubefs/cubefs/util/exporter"
"github.com/cubefs/cubefs/util/log"
"github.com/cubefs/cubefs/util/routinepool"
"github.com/cubefs/cubefs/util/strutil"
)
@ -1856,7 +1855,8 @@ func (s *DataNode) handlePacketToRecoverBackupDataReplica(p *repl.Packet) {
}
}
if rootDir == "" {
log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] dp root not found in dir(%v) .", disk.Path)
err = errors.NewErrorf("dp(%v) root not found in dir(%v)", request.PartitionId, disk.Path)
log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] err %v", err.Error())
return
}
@ -1912,49 +1912,50 @@ func (s *DataNode) handlePacketToRecoverBadDisk(p *repl.Packet) {
return
}
if !disk.startRecover() {
log.LogErrorf("action[handlePacketToRecoverBadDisk] disk(%v) is not found err(%v).", request.DiskPath, err)
err = errors.NewErrorf("disk %v recover is still running")
log.LogErrorf("action[handlePacketToRecoverBadDisk] %v.", err)
return
}
log.LogInfof("action[handlePacketToRecoverBadDisk]req(%v) recover disk %v status %v disk %p",
task.RequestID, disk.Path, disk.recoverStatus, disk)
go func() {
defer disk.stopRecover()
defer func() {
disk.stopRecover()
log.LogInfof("action[handlePacketToRecoverBadDisk]req(%v) recover disk %v async exit status %v", task.RequestID, disk.Path, disk.recoverStatus)
}()
badDpList := disk.GetDiskErrPartitionList()
pool := routinepool.NewRoutinePool(20)
log.LogInfof("action[handlePacketToRecoverBadDisk]req(%v) recover disk %v async enter bad %v disk %p", task.RequestID, disk.Path, len(badDpList), disk)
begin := time.Now()
for _, dpId := range badDpList {
partition := s.space.Partition(dpId)
if partition == nil {
log.LogWarnf("action[handlePacketToRecoverBadDisk] bad dp(%v) not found on disk (%v).", dpId, request.DiskPath)
log.LogWarnf("action[handlePacketToRecoverBadDisk]req(%v) bad dp(%v) not found on disk (%v).", task.RequestID, dpId, request.DiskPath)
continue
}
pool.Submit(func() {
partition.resetDiskErrCnt()
if err = partition.PersistMetadata(); err != nil {
log.LogWarnf("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) persist failed %v.",
partition.partitionID, request.DiskPath, err)
return
}
if err = partition.reload(s.space); err != nil {
log.LogWarnf("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) reload failed %v.",
partition.partitionID, request.DiskPath, err)
return
}
// delete bad io dp
disk.DiskErrPartitionSet.Delete(dpId)
log.LogInfof("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) reload success",
dpId, request.DiskPath)
return
})
partition.resetDiskErrCnt()
if err = partition.PersistMetadata(); err != nil {
log.LogWarnf("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) persist failed %v.",
partition.partitionID, request.DiskPath, err)
continue
}
if err = partition.reload(s.space); err != nil {
log.LogWarnf("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) reload failed %v.",
partition.partitionID, request.DiskPath, err)
continue
}
// delete bad io dp
disk.DiskErrPartitionSet.Delete(dpId)
log.LogInfof("action[handlePacketToRecoverBadDisk] bad dp(%v) on disk (%v) reload success",
task.RequestID, dpId, request.DiskPath)
}
pool.WaitAndClose()
diskErrPartitions := disk.GetDiskErrPartitionList()
if len(diskErrPartitions) != 0 {
err = errors.NewErrorf("disk(%v) has bad dp %v left", request.DiskPath, diskErrPartitions)
return
}
disk.recoverDiskError()
log.LogInfof("action[handlePacketToRecoverBadDisk] recover disk %v success dp(%v) cost %v",
disk.Path, len(badDpList), time.Now().Sub(begin))
log.LogInfof("action[handlePacketToRecoverBadDisk]req(%v) recover disk %v success dp(%v) cost %v",
task.RequestID, disk.Path, len(badDpList), time.Now().Sub(begin))
}()
log.LogInfof("action[handlePacketToRecoverBadDisk] recover bad disk (%v) run async", request.DiskPath)
return

View File

@ -7119,9 +7119,11 @@ func (m *Server) cancelDecommissionDisk(w http.ResponseWriter, r *http.Request)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
// remove from decommissioned disk
m.cluster.deleteAndSyncDecommissionedDisk(dataNode, diskPath)
// remove from decommissioned disk, do not remove bad disk until it is recovered
// dp may allocated on bad disk otherwise
if !dataNode.isBadDisk(diskPath) {
m.cluster.deleteAndSyncDecommissionedDisk(dataNode, diskPath)
}
rstMsg := fmt.Sprintf("cancel decommission disk[%s] successfully ", key)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}
@ -7169,7 +7171,7 @@ func (m *Server) getAllMetaNodes(w http.ResponseWriter, r *http.Request) {
sendOkReply(w, r, newSuccessHTTPReply(metaNodes))
}
func (m *Server) recoverDiskErrorReplica(w http.ResponseWriter, r *http.Request) {
func (m *Server) recoverBackupDataReplica(w http.ResponseWriter, r *http.Request) {
var (
msg string
addr string
@ -7178,9 +7180,9 @@ func (m *Server) recoverDiskErrorReplica(w http.ResponseWriter, r *http.Request)
partitionID uint64
err error
)
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminRecoverDiskErrorReplica))
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminRecoverBackupDataReplica))
defer func() {
doStatAndMetric(proto.AdminRecoverDiskErrorReplica, metric, err, nil)
doStatAndMetric(proto.AdminRecoverBackupDataReplica, metric, err, nil)
}()
if partitionID, addr, err = parseRequestToAddDataReplica(r); err != nil {
@ -7200,13 +7202,13 @@ func (m *Server) recoverDiskErrorReplica(w http.ResponseWriter, r *http.Request)
}
if !proto.IsNormalDp(dp.PartitionType) {
err = fmt.Errorf("action[recoverDiskErrorReplica] [%d] is not normal dp, not support add or delete replica", dp.PartitionID)
err = fmt.Errorf("action[recoverBackupDataReplica] [%d] is not normal dp, not support add or delete replica", dp.PartitionID)
sendErrReply(w, r, newErrHTTPReply(err))
return
}
if dp.ReplicaNum == uint8(len(dp.Replicas)) {
err = fmt.Errorf("action[recoverDiskErrorReplica] [%d] already have %v replicas", dp.PartitionID, dp.ReplicaNum)
err = fmt.Errorf("action[recoverBackupDataReplica] [%d] already have %v replicas", dp.PartitionID, dp.ReplicaNum)
sendErrReply(w, r, newErrHTTPReply(err))
}
@ -7228,27 +7230,27 @@ func (m *Server) recoverDiskErrorReplica(w http.ResponseWriter, r *http.Request)
// restore raft member first
addPeer := proto.Peer{ID: dataNode.ID, Addr: addr}
log.LogInfof("action[recoverDiskErrorReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, dataNode.ID)
log.LogInfof("action[recoverBackupDataReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, dataNode.ID)
if err = m.cluster.addDataPartitionRaftMember(dp, addPeer); err != nil {
log.LogWarnf("action[recoverDiskErrorReplica] dp %v addr %v try add raft member err [%v]", dp.PartitionID, addr, err)
log.LogWarnf("action[recoverBackupDataReplica] dp %v addr %v try add raft member err [%v]", dp.PartitionID, addr, err)
sendErrReply(w, r, newErrHTTPReply(err))
return
}
// find replica disk path
backupInfo, err := dataNode.getBackupDataPartitionInfo(partitionID)
if err != nil {
log.LogWarnf("action[recoverDiskErrorReplica] cannot find backup info for dp %v on dataNode %v", partitionID, dataNode.Addr)
log.LogWarnf("action[recoverBackupDataReplica] cannot find backup info for dp %v on dataNode %v", partitionID, dataNode.Addr)
sendErrReply(w, r, newErrHTTPReply(err))
return
}
err = m.cluster.syncRecoverBackupDataPartitionReplica(addr, backupInfo.Disk, dp)
if err != nil {
log.LogWarnf("action[recoverDiskErrorReplica] dp(%v) recover replica [%v_%v] fail %v", dp.PartitionID, addr, backupInfo.Disk, err)
log.LogWarnf("action[recoverBackupDataReplica] dp(%v) recover replica [%v_%v] fail %v", dp.PartitionID, addr, backupInfo.Disk, err)
sendErrReply(w, r, newErrHTTPReply(err))
return
}
msg = fmt.Sprintf("action[recoverDiskErrorReplica] dp(%v) recover replica [%v_%v] successfully", dp.decommissionInfo(), backupInfo.Addr, backupInfo.Disk)
msg = fmt.Sprintf("action[recoverBackupDataReplica] dp(%v) recover replica [%v_%v] successfully", dp.decommissionInfo(), backupInfo.Addr, backupInfo.Disk)
log.LogInfof("%v", msg)
sendOkReply(w, r, newSuccessHTTPReply(msg))
}
@ -7328,17 +7330,8 @@ func (m *Server) recoverBadDisk(w http.ResponseWriter, r *http.Request) {
return
}
key := fmt.Sprintf("%s_%s", offLineAddr, diskPath)
value, ok := m.cluster.DecommissionDisks.Load(key)
if ok {
disk := value.(*DecommissionDisk)
err = m.cluster.syncDeleteDecommissionDisk(disk)
if err != nil {
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeInternalError, Msg: err.Error()})
return
}
m.cluster.DecommissionDisks.Delete(key)
}
rstMsg := fmt.Sprintf("recover bad disk[%s] successfully ", key)
// do not delete disk from DecommissionDisks, bad disk may be auto decommissioned again if recover is not finished
rstMsg := fmt.Sprintf("recover bad disk[%s] task is submit ", key)
auditlog.LogMasterOp("RecoverBadDisk", rstMsg, nil)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}

View File

@ -4072,6 +4072,7 @@ func (c *Cluster) TryDecommissionDataNode(dataNode *DataNode) {
toBeOffLinePartitions = toBeOffLinePartitions[:dataNode.DecommissionLimit]
}
if len(toBeOffLinePartitions) == 0 {
log.LogWarnf("action[TryDecommissionDataNode]no dp left on dataNode %v, mark decommission success", dataNode.Addr)
dataNode.markDecommissionSuccess(c)
return
}
@ -4313,8 +4314,8 @@ func (c *Cluster) handleDataNodeBadDisk(dataNode *DataNode) {
// decommission failed, but lack replica for disk err dp is already removed
retry := c.RetryDecommissionDisk(dataNode.Addr, disk.DiskPath)
if disk.TotalPartitionCnt == 0 && !retry {
msg := fmt.Sprintf("disk(%v_%v) can be removed", dataNode.Addr, disk.DiskPath)
auditlog.LogMasterOp("DiskDecommission", msg, nil)
// msg := fmt.Sprintf("disk(%v_%v) can be removed", dataNode.Addr, disk.DiskPath)
// auditlog.LogMasterOp("DiskDecommission", msg, nil)
continue
}
var ratio float64

View File

@ -148,7 +148,7 @@ func (dataNode *DataNode) getDisks(c *Cluster) (diskPaths []string) {
return
}
func (dataNode *DataNode) updateBadDisks(latest []string) (ok bool) {
func (dataNode *DataNode) updateBadDisks(latest []string) (ok bool, removed []string) {
sort.Slice(latest, func(i, j int) bool {
return latest[i] < latest[j]
})
@ -157,13 +157,26 @@ func (dataNode *DataNode) updateBadDisks(latest []string) (ok bool) {
dataNode.BadDisks = latest
if len(curr) != len(latest) {
ok = true
return
}
for i := 0; i < len(curr); i++ {
if curr[i] != latest[i] {
ok = true
return
if !ok {
for i := 0; i < len(curr); i++ {
if curr[i] != latest[i] {
ok = true
}
}
}
if ok {
removed = make([]string, 0)
latestMap := make(map[string]bool)
for _, disk := range latest {
latestMap[disk] = true
}
for _, disk := range curr {
if !latestMap[disk] {
removed = append(removed, disk)
}
}
}
return
@ -186,7 +199,7 @@ func (dataNode *DataNode) updateNodeMetric(c *Cluster, resp *proto.DataNodeHeart
dataNode.TotalPartitionSize = resp.TotalPartitionSize
dataNode.AllDisks = resp.AllDisks
updated := dataNode.updateBadDisks(resp.BadDisks)
updated, removedDisks := dataNode.updateBadDisks(resp.BadDisks)
dataNode.BadDiskStats = resp.BadDiskStats
dataNode.DiskStats = resp.DiskStats
dataNode.BackupDataPartitions = resp.BackupDataPartitions
@ -200,6 +213,27 @@ func (dataNode *DataNode) updateNodeMetric(c *Cluster, resp *proto.DataNodeHeart
dataNode.ReportTime = time.Now()
dataNode.isActive = true
if len(removedDisks) != 0 {
log.LogInfof("[updateNodeMetric] dataNode %v removedDisks (%v)", dataNode.Addr, removedDisks)
for _, disk := range removedDisks {
key := fmt.Sprintf("%s_%s", dataNode.Addr, disk)
if value, ok := c.DecommissionDisks.Load(key); ok {
disk := value.(*DecommissionDisk)
if disk.GetDecommissionStatus() == DecommissionCancel {
if err := c.syncDeleteDecommissionDisk(disk); err != nil {
log.LogWarn("[updateNodeMetric] dataNode %v disk (%v) is recovered, but remove failed %v",
dataNode.Addr, key, err)
} else {
c.DecommissionDisks.Delete(key)
// can allocate dp again
c.deleteAndSyncDecommissionedDisk(dataNode, disk.DiskPath)
log.LogInfof("[updateNodeMetric] dataNode %v disk (%v) is recovered", dataNode.Addr, key)
}
}
}
}
}
if updated {
log.LogInfof("[updateNodeMetric] update data node(%v)", dataNode.Addr)
if err := c.syncUpdateDataNode(dataNode); err != nil {
@ -630,3 +664,14 @@ func (dataNode *DataNode) createTaskToQueryBadDiskRecoverProgress(diskPath strin
resp, err = dataNode.TaskManager.syncSendAdminTask(task)
return resp, err
}
func (dataNode *DataNode) isBadDisk(disk string) bool {
dataNode.RLock()
defer dataNode.RUnlock()
for _, entry := range dataNode.BadDisks {
if entry == disk {
return true
}
}
return false
}

View File

@ -1165,6 +1165,7 @@ directly:
}
// initial or failed restart
partition.ResetDecommissionStatus()
partition.DecommissionType = migrateType
partition.SetDecommissionStatus(markDecommission)
partition.DecommissionSrcAddr = srcAddr
partition.DecommissionDstAddr = dstAddr
@ -1284,8 +1285,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
log.LogWarnf("action[decommissionDataPartition] dp [%v] update to prepare failed", partition.PartitionID)
goto errHandler
}
log.LogInfof("action[decommissionDataPartition] dp[%v] from node[%v] to node[%v], raftForce[%v] SingleDecommissionStatus[%v]",
partition.PartitionID, srcAddr, targetAddr, partition.DecommissionRaftForce, partition.GetSpecialReplicaDecommissionStep())
log.LogInfof("action[decommissionDataPartition] dp[%v] start decommission ", partition.decommissionInfo())
// NOTE: delete if not normal data partition or dp is discard
if partition.IsDiscard || !proto.IsNormalDp(partition.PartitionType) {
if vol, ok := c.vols[partition.VolName]; !ok {
@ -2189,14 +2189,19 @@ func (partition *DataPartition) lostLeader(c *Cluster) bool {
}
func (partition *DataPartition) decommissionInfo() string {
var replicas []string
for _, replica := range partition.Replicas {
replicas = append(replicas, replica.Addr)
}
return fmt.Sprintf("vol(%v)_dp(%v)_replicaNum(%v)_src(%v)_dst(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+
"_needRollback(%v)_rollbackTimes(%v)_force(%v)_type(%v)_RestoreReplica(%v)_errMsg(%v)_discard(%v)_term(%v)_addr(%p)",
"_needRollback(%v)_rollbackTimes(%v)_force(%v)_type(%v)_RestoreReplica(%v)_errMsg(%v)_discard(%v)_term(%v)_replica(%v)_addr(%p)",
partition.VolName, partition.PartitionID, partition.ReplicaNum, partition.DecommissionSrcAddr, partition.DecommissionDstAddr,
partition.Hosts, partition.DecommissionRetry, partition.isRecover, GetDecommissionStatusMessage(partition.GetDecommissionStatus()),
GetSpecialDecommissionStatusMessage(partition.GetSpecialReplicaDecommissionStep()), partition.DecommissionNeedRollback,
partition.DecommissionNeedRollbackTimes, partition.DecommissionRaftForce, GetDecommissionTypeMessage(partition.DecommissionType),
GetRestoreReplicaMessage(partition.RestoreReplica), partition.DecommissionErrorMessage, partition.IsDiscard,
partition.DecommissionTerm, partition)
partition.DecommissionTerm, replicas, partition)
}
func (partition *DataPartition) isPerformingDecommission(c *Cluster) bool {

View File

@ -566,8 +566,8 @@ func (m *Server) registerAPIRoutes(router *mux.Router) {
Path(proto.AdminRecoverReplicaMeta).
HandlerFunc(m.recoverReplicaMeta)
router.NewRoute().Methods(http.MethodGet).
Path(proto.AdminRecoverDiskErrorReplica).
HandlerFunc(m.recoverDiskErrorReplica)
Path(proto.AdminRecoverBackupDataReplica).
HandlerFunc(m.recoverBackupDataReplica)
// meta node management APIs
router.NewRoute().Methods(http.MethodGet, http.MethodPost).

View File

@ -47,7 +47,7 @@ const (
AdminQueryDataPartitionDecommissionStatus = "/dataPartition/queryDecommissionStatus"
AdminCheckReplicaMeta = "/dataPartition/checkReplicaMeta"
AdminRecoverReplicaMeta = "/dataPartition/recoverReplicaMeta"
AdminRecoverDiskErrorReplica = "/dataPartition/recoverDiskErrorReplica"
AdminRecoverBackupDataReplica = "/dataPartition/recoverBackupDataReplica"
AdminDeleteDataReplica = "/dataReplica/delete"
AdminAddDataReplica = "/dataReplica/add"
AdminDeleteVol = "/vol/delete"