From bf9fdff28dc3de88ba91d781a185bbb287b8a533 Mon Sep 17 00:00:00 2001 From: chihe Date: Sun, 14 Jul 2024 14:24:37 +0800 Subject: [PATCH] feat(master): save previous decommission error msg when decommissioning another replica Signed-off-by: chihe --- master/api_service.go | 3 +++ master/cluster.go | 8 +++++--- master/data_node.go | 13 +++++++++++++ master/data_partition.go | 12 ++++++++++++ master/disk_manager.go | 34 +++++++++++++++++++++++++++++++++- master/metadata_fsm_op.go | 3 +++ proto/model.go | 1 + 7 files changed, 70 insertions(+), 4 deletions(-) diff --git a/master/api_service.go b/master/api_service.go index b670829a1..3ee91ea53 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -4245,6 +4245,7 @@ func (m *Server) queryDiskDecoProgress(w http.ResponseWriter, r *http.Request) { DecommissionTimes: disk.DecommissionTimes, StatusMessage: GetDecommissionStatusMessage(status), IgnoreDps: disk.IgnoreDecommissionDps, + ResidualDps: disk.residualDecommissionDpsGetAll(), } dps := disk.GetDecommissionFailedDPByTerm(m.cluster) resp.FailedDps = dps @@ -4303,6 +4304,7 @@ func (m *Server) queryAllDecommissionDisk(w http.ResponseWriter, r *http.Request Progress: fmt.Sprintf("%.2f%%", progress*float64(100)), StatusMessage: GetDecommissionStatusMessage(status), IgnoreDps: disk.IgnoreDecommissionDps, + ResidualDps: disk.residualDecommissionDpsGetAll(), FailedDps: disk.GetDecommissionFailedDPByTerm(m.cluster), DecommissionTerm: disk.DecommissionTerm, DecommissionTimes: disk.DecommissionTimes, @@ -5734,6 +5736,7 @@ func (m *Server) queryDataNodeDecoProgress(w http.ResponseWriter, r *http.Reques dps := dn.GetDecommissionFailedDPByTerm(m.cluster) resp.FailedDps = dps resp.IgnoreDps = dn.getIgnoreDecommissionDpList(m.cluster) + resp.ResidualDps = dn.getResidualDecommissionDpList(m.cluster) sendOkReply(w, r, newSuccessHTTPReply(resp)) } diff --git a/master/cluster.go b/master/cluster.go index a11c4a006..b23f06082 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -4201,9 +4201,11 @@ func (c *Cluster) migrateDisk(dataNode *DataNode, diskPath, dstPath string, raft } } else { disk = &DecommissionDisk{ - SrcAddr: nodeAddr, - DiskPath: diskPath, - DiskDisable: diskDisable, + SrcAddr: nodeAddr, + DiskPath: diskPath, + DiskDisable: diskDisable, + IgnoreDecommissionDps: make([]proto.IgnoreDecommissionDP, 0), + ResidualDecommissionDps: make([]proto.IgnoreDecommissionDP, 0), } c.DecommissionDisks.Store(disk.GenerateKey(), disk) } diff --git a/master/data_node.go b/master/data_node.go index abd22f0f7..bb109d0c9 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -697,3 +697,16 @@ func (dataNode *DataNode) getIgnoreDecommissionDpList(c *Cluster) (dps []proto.I } return dps } + +func (dataNode *DataNode) getResidualDecommissionDpList(c *Cluster) (dps []proto.IgnoreDecommissionDP) { + dps = make([]proto.IgnoreDecommissionDP, 0) + for _, disk := range dataNode.DecommissionDiskList { + key := fmt.Sprintf("%s_%s", dataNode.Addr, disk) + // if not found, may already success, so only care running disk + if value, ok := c.DecommissionDisks.Load(key); ok { + dd := value.(*DecommissionDisk) + dps = append(dps, dd.residualDecommissionDpsGetAll()...) + } + } + return dps +} diff --git a/master/data_partition.go b/master/data_partition.go index 8bbe80c6f..9e0fad3df 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -1073,6 +1073,18 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk return errors.NewErrorf("dp[%v] cannot make decommission err:%v", partition.PartitionID, err) } + // if decommission the other replica of this dp, the status of decommission would be overwritten, so save it's error msg + // to last decommission failed disk + if status == DecommissionFail && partition.hasHost(partition.DecommissionSrcAddr) && srcAddr != partition.DecommissionSrcAddr { + key := fmt.Sprintf("%s_%s", partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath) + if value, ok := c.DecommissionDisks.Load(key); ok { + disk := value.(*DecommissionDisk) + if !disk.residualDecommissionDpsHas(partition.PartitionID) { + disk.residualDecommissionDpsSave(partition.PartitionID, partition.DecommissionErrorMessage, c) + } + } + + } // for auto decommission, need raftForce to delete src if no leader if migrateType == AutoDecommission { log.LogDebugf("action[MarkDecommissionStatus] dp[%v] lostLeader %v leader %v interval %v", diff --git a/master/disk_manager.go b/master/disk_manager.go index 1c3d1ae09..affaaa1e4 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -286,9 +286,10 @@ type DecommissionDisk struct { DecommissionDpCount int DiskDisable bool IgnoreDecommissionDps []proto.IgnoreDecommissionDP + ResidualDecommissionDps []proto.IgnoreDecommissionDP Type uint32 DecommissionCompleteTime int64 - UpdateMutex sync.Mutex `json:"-"` + UpdateMutex sync.RWMutex `json:"-"` } func (dd *DecommissionDisk) GenerateKey() string { @@ -536,3 +537,34 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster, ns *nodeSet) (e err = cluster.syncUpdateDecommissionDisk(dd) return err } + +func (dd *DecommissionDisk) residualDecommissionDpsHas(id uint64) bool { + dd.UpdateMutex.RLock() + defer dd.UpdateMutex.RUnlock() + for _, dp := range dd.ResidualDecommissionDps { + if dp.PartitionID == id { + return true + } + } + return false +} + +func (dd *DecommissionDisk) residualDecommissionDpsSave(id uint64, msg string, c *Cluster) { + dd.UpdateMutex.Lock() + defer dd.UpdateMutex.Unlock() + dd.ResidualDecommissionDps = append(dd.ResidualDecommissionDps, proto.IgnoreDecommissionDP{ + PartitionID: id, + ErrMsg: msg, + }) + c.syncUpdateDecommissionDisk(dd) +} + +func (dd *DecommissionDisk) residualDecommissionDpsGetAll() []proto.IgnoreDecommissionDP { + dd.UpdateMutex.RLock() + defer dd.UpdateMutex.RUnlock() + res := make([]proto.IgnoreDecommissionDP, 0) + for _, dp := range dd.ResidualDecommissionDps { + res = append(res, dp) + } + return res +} diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index eeb05d1bb..f553c2c33 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -1725,6 +1725,7 @@ type decommissionDiskValue struct { DecommissionCompleteTime int64 DecommissionLimit int IgnoreDecommissionDps []bsProto.IgnoreDecommissionDP + ResidualDecommissionDps []bsProto.IgnoreDecommissionDP } func newDecommissionDiskValue(disk *DecommissionDisk) *decommissionDiskValue { @@ -1741,6 +1742,7 @@ func newDecommissionDiskValue(disk *DecommissionDisk) *decommissionDiskValue { DecommissionCompleteTime: disk.DecommissionCompleteTime, DecommissionLimit: disk.DecommissionDpCount, IgnoreDecommissionDps: disk.IgnoreDecommissionDps, + ResidualDecommissionDps: disk.ResidualDecommissionDps, } } @@ -1758,6 +1760,7 @@ func (ddv *decommissionDiskValue) Restore() *DecommissionDisk { DecommissionCompleteTime: ddv.DecommissionCompleteTime, DecommissionDpCount: ddv.DecommissionLimit, IgnoreDecommissionDps: ddv.IgnoreDecommissionDps, + ResidualDecommissionDps: ddv.ResidualDecommissionDps, } } diff --git a/proto/model.go b/proto/model.go index c4d59cbde..4510d5fcb 100644 --- a/proto/model.go +++ b/proto/model.go @@ -383,6 +383,7 @@ type DecommissionProgress struct { Progress string FailedDps []FailedDpInfo IgnoreDps []IgnoreDecommissionDP + ResidualDps []IgnoreDecommissionDP } type DiskInfo struct {