From 17b991b3513e3434f8abe6bd8127897f052cc8de Mon Sep 17 00:00:00 2001 From: shuqiang-zheng Date: Mon, 11 Aug 2025 19:25:02 +0800 Subject: [PATCH] feat(master): add status update records to dp decommission process for querying. close: #1000329528 Signed-off-by: shuqiang-zheng --- cli/cmd/const.go | 1 + cli/cmd/datapartition.go | 58 ++++++++++++--- cli/cmd/fmt.go | 15 ++++ master/api_service.go | 26 ++++++- master/cluster.go | 56 +++++++------- master/data_node.go | 2 +- master/data_partition.go | 149 ++++++++++++++++++++++++-------------- master/disk_manager.go | 51 ++++++++----- master/http_server.go | 3 + master/metadata_fsm_op.go | 149 +++++++++++++++++++------------------- master/topology.go | 10 +-- proto/admin_proto.go | 142 ++++++++++++++++++------------------ proto/model.go | 7 ++ sdk/master/api_admin.go | 8 ++ 14 files changed, 413 insertions(+), 264 deletions(-) diff --git a/cli/cmd/const.go b/cli/cmd/const.go index 9ce6ce712..29dc944a3 100644 --- a/cli/cmd/const.go +++ b/cli/cmd/const.go @@ -29,6 +29,7 @@ const ( CliOpDecommission = "decommission" CliOpRecommission = "recommission" CliOpQueryProgress = "query-progress" + CliOpQueryStatusUpdateRecords = "query-status-update-records" CliOpQueryDiskStat = "query-disk-stat" CliOpQueryNodeStat = "query-node-stat" CliOpAbortDecommission = "abort-decommission" diff --git a/cli/cmd/datapartition.go b/cli/cmd/datapartition.go index 17264f26b..b4625fd5a 100644 --- a/cli/cmd/datapartition.go +++ b/cli/cmd/datapartition.go @@ -46,22 +46,24 @@ func newDataPartitionCmd(client *master.MasterClient) *cobra.Command { newDataPartitionResetRestoreStatusCmd(client), newDataPartitionQueryDiskDecommissionInfoStat(client), newDataPartitionQueryDataNodeDecommissionInfoStat(client), + newDataPartitionQueryDecommissionStatusUpdateRecords(client), ) return cmd } const ( - cmdDataPartitionGetShort = "Display detail information of a data partition" - cmdCheckCorruptDataPartitionShort = "Check and list unhealthy data partitions" - cmdDataPartitionDecommissionShort = "Decommission a replication of the data partition to a new address" - cmdDataPartitionReplicateShort = "Add a replication of the data partition on a new address" - cmdDataPartitionDeleteReplicaShort = "Delete a replication of the data partition on a fixed address" - cmdDataPartitionGetDiscardShort = "Display all discard data partitions" - cmdDataPartitionSetDiscardShort = "Set discard flag for data partition" - cmdDataPartitionQueryDecommissionProgressShort = "Query data partition decommission progress" - cmdDataPartitionResetRestoreStatusShort = "Reset data partition restore status" - cmdDataPartitionQueryDiskDecommissionInfoStatShort = "Query data partition disk decommission info stat" - cmdDataPartitionQueryDataNodeDecommissionInfoStatShort = "Query data partition datanode decommission info stat" + cmdDataPartitionGetShort = "Display detail information of a data partition" + cmdCheckCorruptDataPartitionShort = "Check and list unhealthy data partitions" + cmdDataPartitionDecommissionShort = "Decommission a replication of the data partition to a new address" + cmdDataPartitionReplicateShort = "Add a replication of the data partition on a new address" + cmdDataPartitionDeleteReplicaShort = "Delete a replication of the data partition on a fixed address" + cmdDataPartitionGetDiscardShort = "Display all discard data partitions" + cmdDataPartitionSetDiscardShort = "Set discard flag for data partition" + cmdDataPartitionQueryDecommissionProgressShort = "Query data partition decommission progress" + cmdDataPartitionResetRestoreStatusShort = "Reset data partition restore status" + cmdDataPartitionQueryDecommissionStatusUpdateRecordsShort = "Query data partition decommission status update records" + cmdDataPartitionQueryDiskDecommissionInfoStatShort = "Query data partition disk decommission info stat" + cmdDataPartitionQueryDataNodeDecommissionInfoStatShort = "Query data partition datanode decommission info stat" ) func newDataPartitionGetCmd(client *master.MasterClient) *cobra.Command { @@ -680,6 +682,40 @@ func newDataPartitionResetRestoreStatusCmd(client *master.MasterClient) *cobra.C return cmd } +func newDataPartitionQueryDecommissionStatusUpdateRecords(client *master.MasterClient) *cobra.Command { + cmd := &cobra.Command{ + Use: CliOpQueryStatusUpdateRecords + " [DATA PARTITION ID]", + Short: cmdDataPartitionQueryDecommissionStatusUpdateRecordsShort, + Args: cobra.MinimumNArgs(1), + Run: func(cmd *cobra.Command, args []string) { + var ( + err error + dpId uint64 + ) + + defer func() { + errout(err) + }() + + dpId, err = strconv.ParseUint(args[0], 10, 64) + if err != nil { + return + } + + records, err := client.AdminAPI().QueryDataPartitionDecommissionStatusUpdateRecords(dpId) + if err != nil { + return + } + if len(records) == 0 { + stdout("decommission status update records is empty, dp %v may not be in decommissioning\n", dpId) + return + } + stdout("%v", formatDataPartitionDecommissionStatusUpdateRecords(records)) + }, + } + return cmd +} + func newDataPartitionQueryDiskDecommissionInfoStat(client *master.MasterClient) *cobra.Command { cmd := &cobra.Command{ Use: CliOpQueryDiskStat, diff --git a/cli/cmd/fmt.go b/cli/cmd/fmt.go index bea3e719d..d6764a3c6 100644 --- a/cli/cmd/fmt.go +++ b/cli/cmd/fmt.go @@ -1406,6 +1406,21 @@ func formatDataPartitionDecommissionProgress(info *proto.DecommissionDataPartiti return sb.String() } +func formatDataPartitionDecommissionStatusUpdateRecords(records []*proto.DecommissionStatusRecord) string { + sb := strings.Builder{} + if len(records) != 0 { + sb.WriteString("decommission status update records: \n") + for _, record := range records { + sb.WriteString(fmt.Sprintf(" Condition : %v\n", record.Condition)) + sb.WriteString(fmt.Sprintf(" Status : %v\n", record.Status)) + sb.WriteString(fmt.Sprintf(" Time : %v\n", record.Time)) + sb.WriteString(fmt.Sprintf(" ErrMessage : %v\n", record.ErrMessage)) + sb.WriteString("\n") + } + } + return sb.String() +} + func formatDataPartitionDecommissionInfoStat(infos []*proto.DecommissionInfoStat) string { sb := strings.Builder{} if len(infos) != 0 { diff --git a/master/api_service.go b/master/api_service.go index 91c3bcdb0..1dc85ea31 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -1843,7 +1843,7 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) { dp.DecommissionType = ManualAddReplica dp.RecoverStartTime = time.Now() dp.RecoverUpdateTime = time.Now() - dp.SetDecommissionStatus(DecommissionRunning) + dp.SetDecommissionStatus(DecommissionRunning, "manualAddReplica", "") var newReplica *DataReplica if newReplica, err = dp.getReplica(addr); err != nil { @@ -2157,7 +2157,8 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: rstMsg}) return } - err = m.cluster.markDecommissionDataPartition(dp, node, dstNodeSet, raftForce, uint32(decommissionType), weight) + triggerCondition := fmt.Sprintf("manualDecommission_dp(%v)", dp.PartitionID) + err = m.cluster.markDecommissionDataPartition(dp, node, dstNodeSet, raftForce, uint32(decommissionType), weight, triggerCondition) if err != nil { sendErrReply(w, r, newErrHTTPReply(err)) return @@ -2296,6 +2297,27 @@ func (m *Server) resetDataPartitionDecommissionStatus(w http.ResponseWriter, r * sendOkReply(w, r, newSuccessHTTPReply(msg)) } +func (m *Server) queryDataPartitionDecommissionStatusUpdateRecords(w http.ResponseWriter, r *http.Request) { + var ( + dp *DataPartition + partitionID uint64 + err error + records []*proto.DecommissionStatusRecord + ) + + if partitionID, err = parseRequestToLoadDataPartition(r); err != nil { + sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) + return + } + + if dp, err = m.cluster.getDataPartitionByID(partitionID); err != nil { + sendErrReply(w, r, newErrHTTPReply(proto.ErrDataPartitionNotExists)) + return + } + records = dp.cloneDecommissionStatusRecords() + sendOkReply(w, r, newSuccessHTTPReply(records)) +} + func (m *Server) queryDataPartitionDecommissionStatus(w http.ResponseWriter, r *http.Request) { var ( dp *DataPartition diff --git a/master/cluster.go b/master/cluster.go index 486e8e9d5..eaa60aaca 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -516,8 +516,7 @@ func (c *Cluster) scheduleTask() { c.scheduleStartBalanceTask() c.scheduleToUpdateFlashGroupSlots() c.scheduleToCheckDataPartitionRepairingStatus() - c.scheduleToCheckDataPartitionDecommissionDiskRetryMap() - c.scheduleToBalanceDataNode() + c.scheduleToCheckDataPartitionDecommissionInfoRecords() } func (c *Cluster) masterAddr() (addr string) { @@ -2569,7 +2568,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s } // if addDataReplica is success, can add to BadDataPartitionIds dp.SetSpecialReplicaDecommissionStep(SpecialDecommissionWaitAddRes) - dp.SetDecommissionStatus(DecommissionRunning) + dp.SetDecommissionStatus(DecommissionRunning, "decommission_singleDp_addNewReplica", "") dp.isRecover = true dp.Status = proto.ReadOnly dp.RecoverUpdateTime = time.Now() @@ -2587,7 +2586,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s case decommContinue = <-dp.SpecialReplicaDecommissionStop: // if !decommContinue { err = fmt.Errorf("action[decommissionSingleDp] dp %v wait addDataReplica is stopped", dp.PartitionID) - dp.SetDecommissionStatus(DecommissionPause) + dp.SetDecommissionStatus(DecommissionPause, "decommission_singleDp_waitForRepair", err.Error()) log.LogWarnf("action[decommissionSingleDp] dp %v err:%v", dp.PartitionID, err) goto ERR } @@ -2712,7 +2711,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s case decommContinue = <-dp.SpecialReplicaDecommissionStop: if !decommContinue { err = fmt.Errorf("action[decommissionSingleDp] dp %v wait for leader selection is stopped", dp.PartitionID) - dp.SetDecommissionStatus(DecommissionPause) + dp.SetDecommissionStatus(DecommissionPause, "decommission_singleDp_waitForLeader", err.Error()) goto ERR } } @@ -2728,7 +2727,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s goto ERR } dp.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial) - dp.SetDecommissionStatus(DecommissionSuccess) + dp.SetDecommissionStatus(DecommissionSuccess, "decommission_singleDp_deleteOfflineReplica_complete", "") // dp may not add into decommission list when master restart or leader change dp.setRestoreReplicaStop() c.syncUpdateDataPartition(dp) @@ -5299,8 +5298,9 @@ func (c *Cluster) handleDataNodeBadDisk(dataNode *DataNode) { log.LogInfof("[handleDataNodeBadDisk] data node(%v) not found in dp(%v) maybe decommissioned?", dataNode.Addr, dpId) continue } - err = c.markDecommissionDataPartition(dp, dataNode, 0, false, AutoDecommission, highPriorityDecommissionWeight) - if err != nil { + triggerCondition := fmt.Sprintf("autoDecommission_diskErrDp(%v)", dp.PartitionID) + err = c.markDecommissionDataPartition(dp, dataNode, 0, false, AutoDecommission, highPriorityDecommissionWeight, triggerCondition) + if err != nil && !strings.Contains(err.Error(), proto.ErrPerformingDecommission.Error()) { log.LogErrorf("[handleDataNodeBadDisk] failed to decommssion dp(%v) on data node(%v) disk(%v), err(%v)", dataNode.Addr, disk.DiskPath, dp.PartitionID, err) continue } @@ -5332,7 +5332,9 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) { ignorePartitions = make([]*DataPartition, 0) ) defer func() { - auditlog.LogMasterOp("RunningDiskDecommissionIgnoredDps", rstMsg, err) + if len(ignorePartitionIds) != 0 || err != nil { + auditlog.LogMasterOp("RunningDiskDecommissionIgnoredDps", rstMsg, err) + } }() for _, ignoreDecommissionDpInfo := range disk.IgnoreDecommissionDps { @@ -5348,9 +5350,8 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) { disk.SrcAddr, disk.DiskPath, len(ignorePartitionIds), ignorePartitionIds) ignorePartitionIds = ignorePartitionIds[:0] if len(ignorePartitions) == 0 { - log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] no any ignore partitions on disk[%v_%v]", - disk.SrcAddr, disk.DiskPath) rstMsg = fmt.Sprintf("no any ignore partitions on disk[%v]", disk.decommissionInfo()) + log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] %v", rstMsg) return } if node, err = c.dataNode(disk.SrcAddr); err != nil { @@ -5372,7 +5373,7 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) { for _, ignoreDp := range ignorePartitions { triggerCondition := fmt.Sprintf("disk(%v)_%v_dp(%v)", disk.SrcAddr+"_"+disk.DiskPath, disk.Type, ignoreDp.PartitionID) if err = ignoreDp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce, - disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil { + disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns, triggerCondition); err != nil { if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) { c.syncUpdateDataPartition(ignoreDp) // still decommission dp but not involved in the calculation of the decommission progress. @@ -5533,8 +5534,9 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) { ignoreIDs = append(ignoreIDs, dp.PartitionID) continue } + triggerCondition := fmt.Sprintf("disk(%v)_%v_dp(%v)", disk.SrcAddr+"_"+disk.DiskPath, disk.Type, dp.PartitionID) if err = dp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce, - disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil { + disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns, triggerCondition); err != nil { if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) { c.syncUpdateDataPartition(dp) // still decommission dp but not involved in the calculation of the decommission progress. @@ -5575,7 +5577,7 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) { // mark as failed and set decommission src, make sure it can be included in the calculation of progress dp.DecommissionSrcAddr = node.Addr dp.DecommissionSrcDiskPath = disk.DiskPath - dp.markRollbackFailed(false) + dp.markRollbackFailed(false, triggerCondition, err.Error()) dp.DecommissionErrorMessage = err.Error() dp.DecommissionTerm = disk.DecommissionTerm dp.addRetryTimesByDiskPath(dp.DecommissionSrcAddr + "_" + dp.DecommissionSrcDiskPath) @@ -6350,7 +6352,7 @@ func (c *Cluster) rangeAllParitions(f func(d *DataPartition) bool) { } } -func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, dstNodeSetID uint64, raftForce bool, migrateType uint32, weight int) (err error) { +func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, dstNodeSetID uint64, raftForce bool, migrateType uint32, weight int, triggerCondition string) (err error) { addr := src.Addr replica, err := dp.getReplica(addr) if err != nil { @@ -6368,9 +6370,10 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode return } - if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, dstNodeSetID, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns); err != nil { - if !strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) { - dp.markRollbackFailed(false) + if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, dstNodeSetID, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns, triggerCondition); err != nil { + if !strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) && !strings.Contains(err.Error(), proto.ErrPerformingDecommission.Error()) && + !strings.Contains(err.Error(), proto.ErrWaitForAutoAddReplica.Error()) { + dp.markRollbackFailed(false, triggerCondition, err.Error()) dp.DecommissionErrorMessage = err.Error() c.syncUpdateDataPartition(dp) return @@ -6378,9 +6381,9 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode } // TODO: handle error - err = c.syncUpdateDataPartition(dp) - if err != nil { - return + updateErr := c.syncUpdateDataPartition(dp) + if updateErr != nil { + return errors.NewErrorf("dp(%v) mark decommission status failed, err(%v), updateErr(%v)", dp.PartitionID, err, updateErr) } if dp.GetDecommissionStatus() == markDecommission { @@ -6460,25 +6463,28 @@ func (c *Cluster) syncRecoverBackupDataPartitionReplica(host, disk string, dp *D return } -func (c *Cluster) scheduleToCheckDataPartitionDecommissionDiskRetryMap() { +func (c *Cluster) scheduleToCheckDataPartitionDecommissionInfoRecords() { c.runTask(&cTask{ tickTime: time.Second * time.Duration(c.cfg.IntervalToCheckDataPartition), - name: "scheduleToCheckDataPartitionDecommissionDiskRetryMap", + name: "scheduleToCheckDataPartitionDecommissionInfoRecords", function: func() (fin bool) { if c.partition != nil && c.partition.IsRaftLeader() { - c.checkDataPartitionDecommissionDiskRetryMap() + c.checkDataPartitionDecommissionInfoRecords() } return }, }) } -func (c *Cluster) checkDataPartitionDecommissionDiskRetryMap() { +func (c *Cluster) checkDataPartitionDecommissionInfoRecords() { vols := c.allVols() for _, vol := range vols { partitions := vol.dataPartitions.clonePartitions() for _, dp := range partitions { dp.deleteInvalidRetryTimesRecord() + if dp.GetDecommissionStatus() == DecommissionInitial { + dp.clearDecommissionStatusRecords() + } } } } diff --git a/master/data_node.go b/master/data_node.go index ecdd12c5c..00f8672fe 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -678,7 +678,7 @@ func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (remain if dd.GetDecommissionStatus() == markDecommission { remainingDpCnt += dd.GetDecommissionTotalDpCnt(c) } else { - remainingDpCnt += len(partitions) + remainingDpCnt += len(dps) } } } diff --git a/master/data_partition.go b/master/data_partition.go index f6f83c302..7b8ee2e52 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -59,8 +59,9 @@ type DataPartition struct { RdOnly bool addReplicaMutex sync.RWMutex - DecommissionDiskRetryMapMutex sync.RWMutex + DecommissionInfoRecordMutex sync.RWMutex // used for decommissionDiskRetryMap and decommissionStatusUpdateRecords DecommissionDiskRetryMap map[string]int + DecommissionStatusUpdateRecords []*proto.DecommissionStatusRecord DecommissionRetry int DecommissionStatus uint32 DecommissionSrcAddr string @@ -77,17 +78,18 @@ type DataPartition struct { DecommissionWeight int SpecialReplicaDecommissionStop chan bool // used for stop SpecialReplicaDecommissionStep uint32 - IsDiscard bool - VerSeq uint64 - RecoverStartTime time.Time - RecoverUpdateTime time.Time - RecoverLastConsumeTime time.Duration - DecommissionRetryTime time.Time - RepairBlockSize uint64 - DecommissionType uint32 - RestoreReplica uint32 - MediaType uint32 - ForbidWriteOpOfProtoVer0 bool + proto.DecommissionInfoStat + IsDiscard bool + VerSeq uint64 + RecoverStartTime time.Time + RecoverUpdateTime time.Time + RecoverLastConsumeTime time.Duration + DecommissionRetryTime time.Time + RepairBlockSize uint64 + DecommissionType uint32 + RestoreReplica uint32 + MediaType uint32 + ForbidWriteOpOfProtoVer0 bool } func newDataPartition(ID uint64, replicaNum uint8, volName string, volID uint64, @@ -103,6 +105,7 @@ func newDataPartition(ID uint64, replicaNum uint8, volName string, volID uint64, partition.FilesWithMissingReplica = make(map[string]int64) partition.MissingNodes = make(map[string]int64) partition.DecommissionDiskRetryMap = make(map[string]int) + partition.DecommissionStatusUpdateRecords = make([]*proto.DecommissionStatusRecord, 0) partition.Status = proto.ReadOnly partition.VolName = volName @@ -1225,6 +1228,7 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo diskToRepairDpInfo *DiskToDecommissionRepairDpInfo dataNodeParallel uint64 ) + defer c.syncUpdateDataPartition(partition) for _, host := range partition.Hosts { // for AutoAddReplica , firstHost does not need to consider the decommission source address since only adding and not deleting replica @@ -1289,7 +1293,8 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo log.LogInfof("action[AcquireDecommissionFirstHostToken] dp(%v) acquire first host token(%v) success", partition.PartitionID, partition.DecommissionFirstHostDiskTokenKey) return true errHandle: - partition.markRollbackFailed(false) + triggerCondition := fmt.Sprintf("acquireFirsthostToken_firstHost(%v)", firstHost) + partition.markRollbackFailed(false, triggerCondition, err.Error()) partition.DecommissionErrorMessage = err.Error() log.LogWarnf("action[AcquireDecommissionFirstHostToken] clusterID[%v] vol[%v] partitionID[%v]"+ " retry [%v] status [%v] DecommissionDstAddrSpecify [%v] DecommissionDstAddr [%v] DecommissionDstNodeSet [%v] failed", @@ -1298,6 +1303,34 @@ errHandle: return false } +func (partition *DataPartition) recordDecommissionStatus(condition string, errMsg string) { + partition.DecommissionInfoRecordMutex.Lock() + defer partition.DecommissionInfoRecordMutex.Unlock() + record := &proto.DecommissionStatusRecord{ + Condition: condition, + Status: GetDecommissionStatusMessage(partition.DecommissionStatus), + Time: time.Now().Format("2006-01-02 15:04:05"), + ErrMessage: errMsg, + } + partition.DecommissionStatusUpdateRecords = append(partition.DecommissionStatusUpdateRecords, record) +} + +func (partition *DataPartition) cloneDecommissionStatusRecords() []*proto.DecommissionStatusRecord { + partition.DecommissionInfoRecordMutex.RLock() + defer partition.DecommissionInfoRecordMutex.RUnlock() + records := make([]*proto.DecommissionStatusRecord, 0) + records = append(records, partition.DecommissionStatusUpdateRecords...) + return records +} + +func (partition *DataPartition) clearDecommissionStatusRecords() { + partition.DecommissionInfoRecordMutex.Lock() + defer partition.DecommissionInfoRecordMutex.Unlock() + if len(partition.DecommissionStatusUpdateRecords) != 0 { + partition.DecommissionStatusUpdateRecords = make([]*proto.DecommissionStatusRecord, 0) + } +} + func isReplicasContainsHost(replicas []*DataReplica, host string) bool { for _, replica := range replicas { if replica.Addr == host { @@ -1308,7 +1341,7 @@ func isReplicasContainsHost(replicas []*DataReplica, host string) bool { } func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk string, dstNodeSetID uint64, raftForce bool, term uint64, - migrateType uint32, weight int, c *Cluster, ns *nodeSet, + migrateType uint32, weight int, c *Cluster, ns *nodeSet, triggerCondition string, ) (err error) { defer func() { if err != nil { @@ -1547,7 +1580,7 @@ directly: break } partition.DecommissionRetry = 0 - partition.SetDecommissionStatus(markDecommission) + partition.SetDecommissionStatus(markDecommission, triggerCondition, "") // update decommissionTerm for next time query partition.DecommissionTerm = term partition.DecommissionWeight = weight @@ -1577,7 +1610,7 @@ directly: // initial or failed restart partition.ResetDecommissionStatus() partition.DecommissionType = migrateType - partition.SetDecommissionStatus(markDecommission) + partition.SetDecommissionStatus(markDecommission, triggerCondition, "") partition.DecommissionSrcAddr = srcAddr partition.DecommissionDstAddr = dstAddr partition.DecommissionSrcDiskPath = srcDisk @@ -1604,9 +1637,10 @@ directly: return } -func (partition *DataPartition) SetDecommissionStatus(status uint32) { +func (partition *DataPartition) SetDecommissionStatus(status uint32, triggerCondition string, errMsg string) { log.LogDebugf("[SetDecommissionStatus] set dp(%v) decommission status to status(%v)", partition.PartitionID, status) atomic.StoreUint32(&partition.DecommissionStatus, status) + partition.recordDecommissionStatus(triggerCondition, errMsg) } func (partition *DataPartition) SetSpecialReplicaDecommissionStep(step uint32) { @@ -1656,8 +1690,8 @@ func (partition *DataPartition) IsDoingDecommission() bool { } func (partition *DataPartition) cloneDecommissionDiskRetryMap() (result map[string]int) { - partition.DecommissionDiskRetryMapMutex.RLock() - defer partition.DecommissionDiskRetryMapMutex.RUnlock() + partition.DecommissionInfoRecordMutex.RLock() + defer partition.DecommissionInfoRecordMutex.RUnlock() result = make(map[string]int) for disk, retryTimes := range partition.DecommissionDiskRetryMap { result[disk] = retryTimes @@ -1666,8 +1700,8 @@ func (partition *DataPartition) cloneDecommissionDiskRetryMap() (result map[stri } func (partition *DataPartition) addRetryTimesByDiskPath(diskPath string) { - partition.DecommissionDiskRetryMapMutex.Lock() - defer partition.DecommissionDiskRetryMapMutex.Unlock() + partition.DecommissionInfoRecordMutex.Lock() + defer partition.DecommissionInfoRecordMutex.Unlock() if partition.DecommissionDiskRetryMap[diskPath] >= math.MaxInt { partition.DecommissionDiskRetryMap[diskPath] = 0 } else { @@ -1676,29 +1710,29 @@ func (partition *DataPartition) addRetryTimesByDiskPath(diskPath string) { } func (partition *DataPartition) deleteRetryTimesRecordByDiskPath(diskPath string) { - partition.DecommissionDiskRetryMapMutex.Lock() - defer partition.DecommissionDiskRetryMapMutex.Unlock() + partition.DecommissionInfoRecordMutex.Lock() + defer partition.DecommissionInfoRecordMutex.Unlock() delete(partition.DecommissionDiskRetryMap, diskPath) } func (partition *DataPartition) getRetryTimesRecordByDiskPath(diskPath string) (retryTimes int) { - partition.DecommissionDiskRetryMapMutex.RLock() - defer partition.DecommissionDiskRetryMapMutex.RUnlock() + partition.DecommissionInfoRecordMutex.RLock() + defer partition.DecommissionInfoRecordMutex.RUnlock() retryTimes = partition.DecommissionDiskRetryMap[diskPath] return retryTimes } func (partition *DataPartition) deleteInvalidRetryTimesRecord() { - partition.DecommissionDiskRetryMapMutex.RLock() + partition.DecommissionInfoRecordMutex.RLock() if len(partition.DecommissionDiskRetryMap) == 0 { - partition.DecommissionDiskRetryMapMutex.RUnlock() + partition.DecommissionInfoRecordMutex.RUnlock() return } diskRetryMap := make(map[string]int) for disk, retryTimes := range partition.DecommissionDiskRetryMap { diskRetryMap[disk] = retryTimes } - partition.DecommissionDiskRetryMapMutex.RUnlock() + partition.DecommissionInfoRecordMutex.RUnlock() for key := range diskRetryMap { arr := strings.Split(key, "_") if len(arr) == 2 { @@ -1708,9 +1742,9 @@ func (partition *DataPartition) deleteInvalidRetryTimesRecord() { continue } } - partition.DecommissionDiskRetryMapMutex.Lock() + partition.DecommissionInfoRecordMutex.Lock() delete(partition.DecommissionDiskRetryMap, key) - partition.DecommissionDiskRetryMapMutex.Unlock() + partition.DecommissionInfoRecordMutex.Unlock() } } @@ -1733,6 +1767,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { srcAddr = partition.DecommissionSrcAddr targetAddr = partition.DecommissionDstAddr srcReplica *DataReplica + triggerCondition string resetDecommissionDst = true begin = time.Now() finalHosts = make([]string, len(partition.Hosts)) @@ -1741,14 +1776,14 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { if partition.GetDecommissionStatus() == DecommissionInitial { log.LogWarnf("action[decommissionDataPartition] dp [%v] may be cancel", partition.decommissionInfo()) partition.DecommissionErrorMessage = "cancel decommission" - partition.markRollbackFailed(false) + partition.markRollbackFailed(false, "decommission_statusInitial", "cancel decommission") return false } if !c.AutoDecommissionDiskIsEnabled() && partition.DecommissionType == AutoDecommission { log.LogWarnf("action[decommissionDataPartition] dp [%v] decommission is disable", partition.decommissionInfo()) partition.DecommissionErrorMessage = "disable auto " + " decommission" - partition.markRollbackFailed(false) + partition.markRollbackFailed(false, "decommission_autoDecommissionCheck", "disable auto decommission") return false } @@ -1769,11 +1804,11 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { if partition.ReplicaNum == 1 && partition.DecommissionRaftForce { log.LogWarnf("action[decommissionDataPartition] dp [%v] single replica does not support raftForce deletion", partition.decommissionInfo()) partition.DecommissionErrorMessage = "single replica does not support raftForce deletion" - partition.markRollbackFailed(false) + partition.markRollbackFailed(false, "decommission_raftForceCheck", "single replica does not support raftForce deletion") return false } - partition.SetDecommissionStatus(DecommissionPrepare) + partition.SetDecommissionStatus(DecommissionPrepare, "decommission_prepare", "") err = c.syncUpdateDataPartition(partition) if err != nil { log.LogWarnf("action[decommissionDataPartition] dp [%v] update to prepare failed", partition.PartitionID) @@ -1789,7 +1824,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { // log.LogWarnf("[decommissionDataPartition] delete dp(%v) discard(%v)", partition.PartitionID, partition.IsDiscard) // vol.deleteDataPartition(c, partition) // } - partition.SetDecommissionStatus(DecommissionSuccess) + partition.SetDecommissionStatus(DecommissionSuccess, "decommission_discardCheck", "") log.LogWarnf("action[decommissionDataPartition] skip dp(%v) discard(%v)", partition.PartitionID, partition.IsDiscard) return true } @@ -1801,7 +1836,8 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { srcReplica, _ = partition.getReplica(partition.DecommissionSrcAddr) if len(partition.Replicas) == int(partition.ReplicaNum) && srcReplica == nil { - partition.SetDecommissionStatus(DecommissionSuccess) + triggerCondition = fmt.Sprintf("decommission_srcReplica(%v)_hasBeenDeleted", partition.DecommissionSrcAddr) + partition.SetDecommissionStatus(DecommissionSuccess, triggerCondition, "") log.LogWarnf("action[decommissionDataPartition]dp(%v) status(%v) is already decommissioned", partition.PartitionID, partition.Status) return true @@ -1821,7 +1857,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { log.LogWarnf("action[decommissionDataPartition] %s", msg) auditlog.LogMasterOp("DataPartitionDecommission", msg, nil) partition.DecommissionErrorMessage = msg - partition.markRollbackFailed(false) + partition.markRollbackFailed(false, "decommission_raftForceCheck", msg) return false } } @@ -1866,7 +1902,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { newReplica.Status = proto.Recovering // in case heartbeat response is not arrived partition.isRecover = true partition.Status = proto.ReadOnly - partition.SetDecommissionStatus(DecommissionRunning) + partition.SetDecommissionStatus(DecommissionRunning, "decommission_waitForRepair", "") partition.RecoverUpdateTime = time.Now() partition.RecoverStartTime = time.Now() c.putBadDataPartitionIDsByDiskPath(partition.DecommissionSrcDiskPath, partition.DecommissionSrcAddr, partition.PartitionID) @@ -1897,17 +1933,18 @@ errHandler: // if need rollback, set to fail // do not reset DecommissionDstAddr outside the rollback operation, as it may cause rollback failure if partition.DecommissionNeedRollback { - partition.SetDecommissionStatus(DecommissionFail) + partition.SetDecommissionStatus(DecommissionFail, "decommission_needRollBack", err.Error()) } else { // The maximum number of retries for the DP error has been reached, // and a rollback is still required, even if the rollback conditions have not been triggered. if partition.DecommissionRetry >= defaultDecommissionRetryLimit { - partition.markRollbackFailed(true) + triggerCondition = fmt.Sprintf("decommission_retryOverLimit_count(%v)", partition.DecommissionRetry) + partition.markRollbackFailed(true, triggerCondition, err.Error()) } else { // remove dp from BadDataPartitionIDs, preventing errors caused by disk manager not finding the replica - err := c.removeDPFromBadDataPartitionIDs(partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath, partition.PartitionID) - if err != nil { - log.LogWarnf("action[decommissionDataPartition] del dp[%v] from bad dataPartitionIDs failed:%v", partition.PartitionID, err) + removeErr := c.removeDPFromBadDataPartitionIDs(partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath, partition.PartitionID) + if removeErr != nil { + log.LogWarnf("action[decommissionDataPartition] del dp[%v] from bad dataPartitionIDs failed:%v", partition.PartitionID, removeErr) } partition.ReleaseDecommissionToken(c) partition.ReleaseDecommissionFirstHostToken(c) @@ -1916,7 +1953,8 @@ errHandler: partition.DecommissionDstAddr = "" log.LogWarnf("action[decommissionDataPartition] partitionID:%v reset DecommissionDstAddr", partition.PartitionID) } - partition.SetDecommissionStatus(markDecommission) + triggerCondition = fmt.Sprintf("decommission_retry_count(%v)", partition.DecommissionRetry) + partition.SetDecommissionStatus(markDecommission, triggerCondition, err.Error()) } } msg = fmt.Sprintf("clusterID[%v] info[%v] offline failed:%v consume[%v]seconds", @@ -1940,7 +1978,7 @@ func (partition *DataPartition) PauseDecommission(c *Cluster) bool { partition.PartitionID, partition.GetDecommissionStatus()) if status == markDecommission { - partition.SetDecommissionStatus(DecommissionPause) + partition.SetDecommissionStatus(DecommissionPause, "pauseDecommission", "") return true } if partition.isSpecialReplicaCnt() { @@ -1962,7 +2000,7 @@ func (partition *DataPartition) PauseDecommission(c *Cluster) bool { partition.PartitionID, partition.GetDecommissionStatus()) } } - partition.SetDecommissionStatus(DecommissionPause) + partition.SetDecommissionStatus(DecommissionPause, "pauseDecommission", "") partition.isRecover = false return true } @@ -1980,13 +2018,14 @@ func (partition *DataPartition) ResetDecommissionStatus() { partition.DecommissionDstNodeSet = 0 partition.DecommissionNeedRollback = false atomic.StoreUint32(&partition.DecommissionNeedRollbackTimes, 0) - partition.SetDecommissionStatus(DecommissionInitial) + partition.SetDecommissionStatus(DecommissionInitial, "resetDecommissionStatus", "") partition.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial) partition.DecommissionErrorMessage = "" partition.DecommissionType = InitialDecommission partition.RecoverStartTime = time.Time{} partition.RecoverUpdateTime = time.Time{} partition.DecommissionRetryTime = time.Time{} + partition.clearDecommissionStatusRecords() } func (partition *DataPartition) resetRestoreMeta(expected uint32) (ok bool) { @@ -2030,7 +2069,7 @@ func (partition *DataPartition) rollback(c *Cluster) { partition.isRecover = false partition.DecommissionNeedRollback = false partition.DecommissionErrorMessage = "" - partition.SetDecommissionStatus(markDecommission) + partition.SetDecommissionStatus(markDecommission, "rollback_complete", "") partition.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial) // specify dst addr do not need rollback // keep DecommissionSrcAddr to prevent allocate DecommissionSrcAddr data node during acquire token @@ -2413,7 +2452,8 @@ errHandler: partition.DecommissionRetry++ partition.DecommissionRetryTime = time.Now() if partition.DecommissionRetry >= defaultDecommissionRetryLimit { - partition.markRollbackFailed(false) + triggerCondition := "acquireNsDecommissionToken" + partition.markRollbackFailed(false, triggerCondition, err.Error()) } partition.DecommissionErrorMessage = err.Error() log.LogWarnf("action[TryAcquireDecommissionToken] clusterID[%v] vol[%v] partitionID[%v]"+ @@ -2502,8 +2542,8 @@ func (partition *DataPartition) needRollback(c *Cluster) bool { return true } -func (partition *DataPartition) markRollbackFailed(needRollback bool) { - partition.SetDecommissionStatus(DecommissionFail) +func (partition *DataPartition) markRollbackFailed(needRollback bool, triggerCondition string, errMsg string) { + partition.SetDecommissionStatus(DecommissionFail, triggerCondition, errMsg) partition.DecommissionNeedRollbackTimes = defaultDecommissionRollbackLimit partition.DecommissionNeedRollback = needRollback } @@ -2875,7 +2915,8 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) { partition.PartitionID, addr) return nil } - err = c.markDecommissionDataPartition(partition, node, 0, false, AutoAddReplica, highPriorityDecommissionWeight) + triggerCondition := fmt.Sprintf("autoAddReplica_dp(%v)", partition.PartitionID) + err = c.markDecommissionDataPartition(partition, node, 0, false, AutoAddReplica, highPriorityDecommissionWeight, triggerCondition) auditMsg = fmt.Sprintf("dp(%v) ReplicaNum %v hostsNum %v auto add replica", partition.PartitionID, partition.ReplicaNum, len(partition.Hosts)) log.LogDebugf("action[checkReplicaMeta]%v: err %v", auditMsg, err) @@ -3000,11 +3041,11 @@ func (partition *DataPartition) removeHostByForce(c *Cluster, peerAddr string) { } } -func (partition *DataPartition) resetForManualAddReplica() { +func (partition *DataPartition) resetForManualAddReplica(triggerCondition string, errMsg string) { partition.DecommissionDstAddr = "" partition.DecommissionType = InitialDecommission partition.isRecover = false - partition.SetDecommissionStatus(DecommissionInitial) + partition.SetDecommissionStatus(DecommissionInitial, triggerCondition, errMsg) partition.setRestoreReplicaStop() } diff --git a/master/disk_manager.go b/master/disk_manager.go index e1bc1d65e..806183fe5 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -77,12 +77,12 @@ func (c *Cluster) checkDiskRecoveryProgress() { log.LogInfof("action[checkDiskRecoveryProgress] dp %v isSpec %v replicas %v conf replicas num %v status(%v)", partition.decommissionInfo(), partition.isSpecialReplicaCnt(), len(partition.Replicas), int(partition.ReplicaNum), partition.GetDecommissionStatus()) if len(partition.Replicas) == 0 { - partition.SetDecommissionStatus(DecommissionSuccess) + partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress_dpMaybeDeleted", "") log.LogWarnf("action[checkDiskRecoveryProgress] dp %v maybe deleted", partition.PartitionID) continue } if partition.IsDiscard { - partition.SetDecommissionStatus(DecommissionSuccess) + partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress_discardCheck", "") log.LogWarnf("[checkDiskRecoveryProgress] dp(%v) is discard, decommission successfully", partition.PartitionID) continue } @@ -97,15 +97,16 @@ func (c *Cluster) checkDiskRecoveryProgress() { newReplica, _ := partition.getReplica(partition.DecommissionDstAddr) if newReplica == nil { + errMsg := fmt.Sprintf("Decommission target node %v not found", partition.DecommissionDstAddr) log.LogWarnf("action[checkDiskRecoveryProgress] dp %v cannot find replica %v", partition.PartitionID, partition.DecommissionDstAddr) if partition.DecommissionType == ManualAddReplica { - partition.resetForManualAddReplica() + partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg) } else { partition.DecommissionNeedRollback = true - partition.SetDecommissionStatus(DecommissionFail) + partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg) } - partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v not found", partition.DecommissionDstAddr) + partition.DecommissionErrorMessage = errMsg partition.RLock() err = c.syncUpdateDataPartition(partition) if err != nil { @@ -122,18 +123,20 @@ func (c *Cluster) checkDiskRecoveryProgress() { duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0)) diskErrReplicas := partition.getAllDiskErrorReplica() if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) || math.Abs(duration.Minutes()) > 10 { - if partition.DecommissionType == ManualAddReplica { - partition.resetForManualAddReplica() - } else { - partition.markRollbackFailed(true) - } + var errMsg string if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) { - partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+ + errMsg = fmt.Sprintf("Decommission target node %v cannot finish recover"+ " for host[0] %v is unavailable", partition.DecommissionDstAddr, partition.Hosts[0]) } else { - partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+ + errMsg = fmt.Sprintf("Decommission target node %v cannot finish recover"+ " for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr) } + if partition.DecommissionType == ManualAddReplica { + partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg) + } else { + partition.markRollbackFailed(true, "checkDiskRecoveryProgress", errMsg) + } + partition.DecommissionErrorMessage = errMsg Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v", c.Name, partitionID, partition.DecommissionErrorMessage)) partition.RLock() @@ -144,13 +147,14 @@ func (c *Cluster) checkDiskRecoveryProgress() { partition.RUnlock() continue } else if time.Since(partition.RecoverUpdateTime) > c.GetDecommissionDataPartitionRecoverTimeOut() { + errMsg := fmt.Sprintf("Decommission target node %v repair timeout", partition.DecommissionDstAddr) if partition.DecommissionType == ManualAddReplica { - partition.resetForManualAddReplica() + partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg) } else { partition.DecommissionNeedRollback = true - partition.SetDecommissionStatus(DecommissionFail) + partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg) } - partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v repair timeout", partition.DecommissionDstAddr) + partition.DecommissionErrorMessage = errMsg Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v_%v recovered timeout,recoverUpdateTime %s", c.Name, partitionID, newReplica.Addr, newReplica.DiskPath, time.Since(partition.RecoverUpdateTime))) partition.RLock() @@ -165,8 +169,10 @@ func (c *Cluster) checkDiskRecoveryProgress() { newBadDpIds = append(newBadDpIds, partitionID) } else { if partition.DecommissionType == ManualAddReplica { + var errMsg string if newReplica.isUnavailable() { - partition.DecommissionErrorMessage = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr) + errMsg = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr) + partition.DecommissionErrorMessage = errMsg Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered failed", c.Name, partitionID, partition.DecommissionDstAddr)) } else { @@ -174,7 +180,11 @@ func (c *Cluster) checkDiskRecoveryProgress() { Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered success", c.Name, partitionID, partition.DecommissionDstAddr)) } - partition.resetForManualAddReplica() + partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg) + if errMsg == "" { + partition.clearDecommissionStatusRecords() + } + log.LogInfof("[checkDiskRecoveryProgress] dp(%v) manual add new replica addr %v status(%v)", partitionID, newReplica.Addr, newReplica.Status) partition.RLock() @@ -192,14 +202,15 @@ func (c *Cluster) checkDiskRecoveryProgress() { } // do not add to BadDataPartitionIds if newReplica.isUnavailable() { + errMsg := fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr) partition.DecommissionNeedRollback = true - partition.SetDecommissionStatus(DecommissionFail) - partition.DecommissionErrorMessage = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr) + partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg) + partition.DecommissionErrorMessage = errMsg Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered failed", c.Name, partitionID, partition.DecommissionDstAddr)) } else { partition.DecommissionErrorMessage = "" - partition.SetDecommissionStatus(DecommissionSuccess) // can be readonly or readwrite + partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress", "") // can be readonly or readwrite Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] "+ "replica %v has recovered success,cost(%v)", c.Name, partitionID, partition.DecommissionDstAddr, time.Since(partition.RecoverStartTime).String())) diff --git a/master/http_server.go b/master/http_server.go index 1ca92d2c8..e5bbec403 100644 --- a/master/http_server.go +++ b/master/http_server.go @@ -634,6 +634,9 @@ func (m *Server) registerAPIRoutes(router *mux.Router) { router.NewRoute().Methods(http.MethodGet). Path(proto.AdminQueryDataPartitionDecommissionStatus). HandlerFunc(m.queryDataPartitionDecommissionStatus) + router.NewRoute().Methods(http.MethodGet). + Path(proto.AdminQueryDataPartitionDecommissionStatusUpdateRecords). + HandlerFunc(m.queryDataPartitionDecommissionStatusUpdateRecords) router.NewRoute().Methods(http.MethodGet). Path(proto.AdminCheckReplicaMeta). HandlerFunc(m.checkReplicaMeta) diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index 641bd6fa1..1ddc0fe83 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -179,42 +179,43 @@ func newMetaPartitionValue(mp *MetaPartition) (mpv *metaPartitionValue) { } type dataPartitionValue struct { - PartitionID uint64 - ReplicaNum uint8 - Hosts string - Peers []proto.Peer - Status int8 - VolID uint64 - VolName string - OfflinePeerID uint64 - Replicas []*replicaValue - IsRecover bool - PartitionType int - RdOnly bool - IsDiscard bool - DecommissionDiskRetryMap map[string]int - DecommissionRetry int - DecommissionStatus uint32 - DecommissionSrcAddr string - DecommissionDstAddr string - DecommissionRaftForce bool - DecommissionSrcDiskPath string - DecommissionTerm uint64 - DecommissionWeight int - SpecialReplicaDecommissionStep uint32 - DecommissionDstAddrSpecify bool - DecommissionDstNodeSet uint64 - DecommissionNeedRollback bool - RecoverStartTime int64 - RecoverUpdateTime int64 - RecoverLastConsumeTime float64 - DecommissionRetryTime int64 - Forbidden bool - DecommissionErrorMessage string - DecommissionNeedRollbackTimes uint32 - DecommissionType uint32 - RestoreReplica uint32 - MediaType uint32 + PartitionID uint64 + ReplicaNum uint8 + Hosts string + Peers []proto.Peer + Status int8 + VolID uint64 + VolName string + OfflinePeerID uint64 + Replicas []*replicaValue + IsRecover bool + PartitionType int + RdOnly bool + IsDiscard bool + DecommissionDiskRetryMap map[string]int + DecommissionStatusUpdateRecords []*proto.DecommissionStatusRecord + DecommissionRetry int + DecommissionStatus uint32 + DecommissionSrcAddr string + DecommissionDstAddr string + DecommissionRaftForce bool + DecommissionSrcDiskPath string + DecommissionTerm uint64 + DecommissionWeight int + SpecialReplicaDecommissionStep uint32 + DecommissionDstAddrSpecify bool + DecommissionDstNodeSet uint64 + DecommissionNeedRollback bool + RecoverStartTime int64 + RecoverUpdateTime int64 + RecoverLastConsumeTime float64 + DecommissionRetryTime int64 + Forbidden bool + DecommissionErrorMessage string + DecommissionNeedRollbackTimes uint32 + DecommissionType uint32 + RestoreReplica uint32 + MediaType uint32 } func (dpv *dataPartitionValue) Restore(c *Cluster) (dp *DataPartition) { @@ -268,6 +269,7 @@ func (dpv *dataPartitionValue) Restore(c *Cluster) (dp *DataPartition) { for disk, retryTimes := range dpv.DecommissionDiskRetryMap { dp.DecommissionDiskRetryMap[disk] = retryTimes } + dp.DecommissionStatusUpdateRecords = append(dp.DecommissionStatusUpdateRecords, dpv.DecommissionStatusUpdateRecords...) return dp } @@ -278,50 +280,47 @@ type replicaValue struct { func newDataPartitionValue(dp *DataPartition) (dpv *dataPartitionValue) { dpv = &dataPartitionValue{ - PartitionID: dp.PartitionID, - ReplicaNum: dp.ReplicaNum, - Hosts: dp.hostsToString(), - Peers: dp.Peers, - Status: dp.Status, - VolID: dp.VolID, - VolName: dp.VolName, - OfflinePeerID: dp.OfflinePeerID, - Replicas: make([]*replicaValue, 0), - IsRecover: dp.isRecover, - PartitionType: dp.PartitionType, - RdOnly: dp.RdOnly, - IsDiscard: dp.IsDiscard, - DecommissionDiskRetryMap: make(map[string]int), - DecommissionRetry: dp.DecommissionRetry, - DecommissionStatus: atomic.LoadUint32(&dp.DecommissionStatus), - DecommissionSrcAddr: dp.DecommissionSrcAddr, - DecommissionDstAddr: dp.DecommissionDstAddr, - DecommissionRaftForce: dp.DecommissionRaftForce, - DecommissionSrcDiskPath: dp.DecommissionSrcDiskPath, - DecommissionTerm: dp.DecommissionTerm, - DecommissionWeight: dp.DecommissionWeight, - SpecialReplicaDecommissionStep: dp.SpecialReplicaDecommissionStep, - DecommissionDstAddrSpecify: dp.DecommissionDstAddrSpecify, - DecommissionDstNodeSet: dp.DecommissionDstNodeSet, - DecommissionNeedRollback: dp.DecommissionNeedRollback, - RecoverStartTime: dp.RecoverStartTime.Unix(), - RecoverUpdateTime: dp.RecoverUpdateTime.Unix(), - RecoverLastConsumeTime: dp.RecoverLastConsumeTime.Seconds(), - DecommissionRetryTime: dp.DecommissionRetryTime.Unix(), - DecommissionErrorMessage: dp.DecommissionErrorMessage, - DecommissionNeedRollbackTimes: dp.DecommissionNeedRollbackTimes, - DecommissionType: dp.DecommissionType, - RestoreReplica: atomic.LoadUint32(&dp.RestoreReplica), - MediaType: dp.MediaType, + PartitionID: dp.PartitionID, + ReplicaNum: dp.ReplicaNum, + Hosts: dp.hostsToString(), + Peers: dp.Peers, + Status: dp.Status, + VolID: dp.VolID, + VolName: dp.VolName, + OfflinePeerID: dp.OfflinePeerID, + Replicas: make([]*replicaValue, 0), + IsRecover: dp.isRecover, + PartitionType: dp.PartitionType, + RdOnly: dp.RdOnly, + IsDiscard: dp.IsDiscard, + DecommissionDiskRetryMap: dp.cloneDecommissionDiskRetryMap(), + DecommissionStatusUpdateRecords: dp.cloneDecommissionStatusRecords(), + DecommissionRetry: dp.DecommissionRetry, + DecommissionStatus: atomic.LoadUint32(&dp.DecommissionStatus), + DecommissionSrcAddr: dp.DecommissionSrcAddr, + DecommissionDstAddr: dp.DecommissionDstAddr, + DecommissionRaftForce: dp.DecommissionRaftForce, + DecommissionSrcDiskPath: dp.DecommissionSrcDiskPath, + DecommissionTerm: dp.DecommissionTerm, + DecommissionWeight: dp.DecommissionWeight, + SpecialReplicaDecommissionStep: dp.SpecialReplicaDecommissionStep, + DecommissionDstAddrSpecify: dp.DecommissionDstAddrSpecify, + DecommissionDstNodeSet: dp.DecommissionDstNodeSet, + DecommissionNeedRollback: dp.DecommissionNeedRollback, + RecoverStartTime: dp.RecoverStartTime.Unix(), + RecoverUpdateTime: dp.RecoverUpdateTime.Unix(), + RecoverLastConsumeTime: dp.RecoverLastConsumeTime.Seconds(), + DecommissionRetryTime: dp.DecommissionRetryTime.Unix(), + DecommissionErrorMessage: dp.DecommissionErrorMessage, + DecommissionNeedRollbackTimes: dp.DecommissionNeedRollbackTimes, + DecommissionType: dp.DecommissionType, + RestoreReplica: atomic.LoadUint32(&dp.RestoreReplica), + MediaType: dp.MediaType, } for _, replica := range dp.Replicas { rv := &replicaValue{Addr: replica.Addr, DiskPath: replica.DiskPath} dpv.Replicas = append(dpv.Replicas, rv) } - retryTimesMap := dp.cloneDecommissionDiskRetryMap() - for disk, retryTimes := range retryTimesMap { - dpv.DecommissionDiskRetryMap[disk] = retryTimes - } return } diff --git a/master/topology.go b/master/topology.go index 2b98f9a52..d565150f2 100644 --- a/master/topology.go +++ b/master/topology.go @@ -2197,7 +2197,7 @@ func (l *DecommissionDataPartitionList) Put(id uint64, value *DataPartition, c * } // prepare status reset to mark status to retry again if value.GetDecommissionStatus() == DecommissionPrepare { - value.SetDecommissionStatus(markDecommission) + value.SetDecommissionStatus(markDecommission, "leaderChange_updatePrepareToMark", "") } l.mu.Lock() if _, ok := l.cacheMap[value.PartitionID]; ok { @@ -2223,7 +2223,7 @@ func (l *DecommissionDataPartitionList) Put(id uint64, value *DataPartition, c * // restore special replica decommission progress if value.isSpecialReplicaCnt() && value.GetDecommissionStatus() == DecommissionRunning && !value.DecommissionRaftForce { - value.SetDecommissionStatus(markDecommission) + value.SetDecommissionStatus(markDecommission, "leaderChange_updateRunningToMark", "") value.isRecover = false // can pass decommission validate check log.LogInfof("action[DecommissionDataPartitionListPut] ns[%v] dp[%v] set status from DecommissionRunning to markDecommission", id, value.PartitionID) @@ -2323,16 +2323,16 @@ func (l *DecommissionDataPartitionList) startTraverse() { func updateDecommissionWeight(dps []*DataPartition, c *Cluster) { for _, dp := range dps { if dp.IsDiscard { - dp.SetDecommissionStatus(DecommissionSuccess) + dp.SetDecommissionStatus(DecommissionSuccess, "traverDecommissionList_updateDecommissionWeight_discardCheck", "") log.LogWarnf("action[DecommissionListTraverse] skip dp(%v) discard(%v)", dp.PartitionID, dp.IsDiscard) continue } diskErrReplicaNum := dp.getReplicaDiskErrorNum() if diskErrReplicaNum == dp.ReplicaNum || diskErrReplicaNum == uint8(len(dp.Peers)) { - log.LogWarnf("action[DecommissionListTraverse] dp[%v] all live replica is unavaliable", dp.decommissionInfo()) + log.LogWarnf("action[DecommissionListTraverse] dp[%v] all live replica is unavailable", dp.decommissionInfo()) err := proto.ErrAllReplicaUnavailable dp.DecommissionErrorMessage = err.Error() - dp.markRollbackFailed(false) + dp.markRollbackFailed(false, "traverDecommissionList_updateDecommissionWeight_diskErrReplicaNumCheck", err.Error()) continue } if dp.DecommissionType == AutoDecommission && dp.IsMarkDecommission() { diff --git a/proto/admin_proto.go b/proto/admin_proto.go index 75fc4476c..18870d462 100644 --- a/proto/admin_proto.go +++ b/proto/admin_proto.go @@ -34,77 +34,77 @@ type ContextUserKey string // api const ( // Admin APIs - AdminGetMasterApiList = "/admin/getMasterApiList" - AdminSetApiQpsLimit = "/admin/setApiQpsLimit" - AdminGetApiQpsLimit = "/admin/getApiQpsLimit" - AdminRemoveApiQpsLimit = "/admin/rmApiQpsLimit" - AdminGetCluster = "/admin/getCluster" - AdminSetClusterInfo = "/admin/setClusterInfo" - AdminGetMonitorPushAddr = "/admin/getMonitorPushAddr" - AdminGetClusterDataNodes = "/admin/cluster/getAllDataNodes" - AdminGetClusterMetaNodes = "/admin/cluster/getAllMetaNodes" - AdminGetDataPartition = "/dataPartition/get" - AdminLoadDataPartition = "/dataPartition/load" - AdminCreateDataPartition = "/dataPartition/create" - AdminCreatePreLoadDataPartition = "/dataPartition/createPreLoad" - AdminDecommissionDataPartition = "/dataPartition/decommission" - AdminDiagnoseDataPartition = "/dataPartition/diagnose" - AdminResetDataPartitionDecommissionStatus = "/dataPartition/resetDecommissionStatus" - AdminQueryDataPartitionDecommissionStatus = "/dataPartition/queryDecommissionStatus" - AdminCheckReplicaMeta = "/dataPartition/checkReplicaMeta" - AdminRecoverReplicaMeta = "/dataPartition/recoverReplicaMeta" - AdminRecoverBackupDataReplica = "/dataPartition/recoverBackupDataReplica" - AdminDeleteDataReplica = "/dataReplica/delete" - AdminAddDataReplica = "/dataReplica/add" - AdminDeleteVol = "/vol/delete" - AdminUpdateVol = "/vol/update" - AdminVolShrink = "/vol/shrink" - AdminVolExpand = "/vol/expand" - AdminVolForbidden = "/vol/forbidden" - AdminVolEnableAuditLog = "/vol/auditlog" - AdminVolSetDpRepairBlockSize = "/vol/setDpRepairBlockSize" - AdminCreateVol = "/admin/createVol" - AdminGetVol = "/admin/getVol" - AdminClusterFreeze = "/cluster/freeze" - AdminClusterForbidMpDecommission = "/cluster/forbidMetaPartitionDecommission" - AdminClusterStat = "/cluster/stat" - AdminSetCheckDataReplicasEnable = "/cluster/setCheckDataReplicasEnable" - AdminGetIP = "/admin/getIp" - AdminCreateMetaPartition = "/metaPartition/create" - AdminSetMetaNodeThreshold = "/threshold/set" - AdminSetMasterVolDeletionDelayTime = "/volDeletionDelayTime/set" - AdminSetMetaNodeGOGC = "/metaNodeGOGC/set" - AdminSetDataNodeGOGC = "/dataNodeGOGC/set" - AdminListVols = "/vol/list" - AdminSetNodeInfo = "/admin/setNodeInfo" - AdminGetNodeInfo = "/admin/getNodeInfo" - AdminGetAllNodeSetGrpInfo = "/admin/getDomainInfo" - AdminGetNodeSetGrpInfo = "/admin/getDomainNodeSetGrpInfo" - AdminGetIsDomainOn = "/admin/getIsDomainOn" - AdminUpdateNodeSetCapcity = "/admin/updateNodeSetCapcity" - AdminUpdateNodeSetId = "/admin/updateNodeSetId" - AdminUpdateNodeSetNodeSelector = "/admin/updateNodeSetNodeSelector" - AdminUpdateDomainDataUseRatio = "/admin/updateDomainDataRatio" - AdminUpdateZoneExcludeRatio = "/admin/updateZoneExcludeRatio" - AdminSetNodeRdOnly = "/admin/setNodeRdOnly" - AdminSetDpRdOnly = "/admin/setDpRdOnly" - AdminSetConfig = "/admin/setConfig" - AdminGetConfig = "/admin/getConfig" - AdminDataPartitionChangeLeader = "/dataPartition/changeleader" - AdminChangeMasterLeader = "/master/changeleader" - AdminOpFollowerPartitionsRead = "/master/opFollowerPartitionRead" - AdminUpdateDecommissionFirstHostDiskParallelLimit = "/admin/updateDecommissionFirstHostDiskParallelLimit" - AdminQueryDecommissionFirstHostDiskParallelLimit = "/admin/queryDecommissionFirstHostDiskParallelLimit" - AdminUpdateDecommissionFirstHostParallelLimit = "/admin/updateDecommissionFirstHostParallelLimit" - AdminQueryDecommissionFirstHostParallelLimit = "/admin/queryDecommissionFirstHostParallelLimit" - AdminQueryDecommissionFirstHostParallelInfo = "/admin/queryDecommissionFirstHostParallelInfo" - AdminUpdateDecommissionLimit = "/admin/updateDecommissionLimit" - AdminQueryDecommissionLimit = "/admin/queryDecommissionLimit" - AdminQueryDecommissionFailedDisk = "/admin/queryDecommissionFailedDisk" - AdminAbortDecommissionDisk = "/admin/abortDecommissionDisk" - AdminResetDataPartitionRestoreStatus = "/admin/resetDataPartitionRestoreStatus" - AdminGetOpLog = "/admin/getOpLog" - AdminGetRemoteCacheConfig = "/admin/getRemoteCacheConfig" + AdminGetMasterApiList = "/admin/getMasterApiList" + AdminSetApiQpsLimit = "/admin/setApiQpsLimit" + AdminGetApiQpsLimit = "/admin/getApiQpsLimit" + AdminRemoveApiQpsLimit = "/admin/rmApiQpsLimit" + AdminGetCluster = "/admin/getCluster" + AdminSetClusterInfo = "/admin/setClusterInfo" + AdminGetMonitorPushAddr = "/admin/getMonitorPushAddr" + AdminGetClusterDataNodes = "/admin/cluster/getAllDataNodes" + AdminGetClusterMetaNodes = "/admin/cluster/getAllMetaNodes" + AdminGetDataPartition = "/dataPartition/get" + AdminLoadDataPartition = "/dataPartition/load" + AdminCreateDataPartition = "/dataPartition/create" + AdminDecommissionDataPartition = "/dataPartition/decommission" + AdminDiagnoseDataPartition = "/dataPartition/diagnose" + AdminResetDataPartitionDecommissionStatus = "/dataPartition/resetDecommissionStatus" + AdminQueryDataPartitionDecommissionStatus = "/dataPartition/queryDecommissionStatus" + AdminQueryDataPartitionDecommissionStatusUpdateRecords = "/dataPartition/queryDecommissionStatusUpdateRecords" + AdminCheckReplicaMeta = "/dataPartition/checkReplicaMeta" + AdminRecoverReplicaMeta = "/dataPartition/recoverReplicaMeta" + AdminRecoverBackupDataReplica = "/dataPartition/recoverBackupDataReplica" + AdminDeleteDataReplica = "/dataReplica/delete" + AdminAddDataReplica = "/dataReplica/add" + AdminDeleteVol = "/vol/delete" + AdminUpdateVol = "/vol/update" + AdminVolShrink = "/vol/shrink" + AdminVolExpand = "/vol/expand" + AdminVolForbidden = "/vol/forbidden" + AdminVolEnableAuditLog = "/vol/auditlog" + AdminVolSetDpRepairBlockSize = "/vol/setDpRepairBlockSize" + AdminCreateVol = "/admin/createVol" + AdminGetVol = "/admin/getVol" + AdminClusterFreeze = "/cluster/freeze" + AdminClusterForbidMpDecommission = "/cluster/forbidMetaPartitionDecommission" + AdminClusterStat = "/cluster/stat" + AdminSetCheckDataReplicasEnable = "/cluster/setCheckDataReplicasEnable" + AdminGetIP = "/admin/getIp" + AdminCreateMetaPartition = "/metaPartition/create" + AdminSetMetaNodeThreshold = "/threshold/set" + AdminSetMasterVolDeletionDelayTime = "/volDeletionDelayTime/set" + AdminSetMetaNodeGOGC = "/metaNodeGOGC/set" + AdminSetDataNodeGOGC = "/dataNodeGOGC/set" + AdminListVols = "/vol/list" + AdminSetNodeInfo = "/admin/setNodeInfo" + AdminGetNodeInfo = "/admin/getNodeInfo" + AdminGetAllNodeSetGrpInfo = "/admin/getDomainInfo" + AdminGetNodeSetGrpInfo = "/admin/getDomainNodeSetGrpInfo" + AdminGetIsDomainOn = "/admin/getIsDomainOn" + AdminUpdateNodeSetCapcity = "/admin/updateNodeSetCapcity" + AdminUpdateNodeSetId = "/admin/updateNodeSetId" + AdminUpdateNodeSetNodeSelector = "/admin/updateNodeSetNodeSelector" + AdminUpdateDomainDataUseRatio = "/admin/updateDomainDataRatio" + AdminUpdateZoneExcludeRatio = "/admin/updateZoneExcludeRatio" + AdminSetNodeRdOnly = "/admin/setNodeRdOnly" + AdminSetDpRdOnly = "/admin/setDpRdOnly" + AdminSetConfig = "/admin/setConfig" + AdminGetConfig = "/admin/getConfig" + AdminDataPartitionChangeLeader = "/dataPartition/changeleader" + AdminChangeMasterLeader = "/master/changeleader" + AdminOpFollowerPartitionsRead = "/master/opFollowerPartitionRead" + AdminUpdateDecommissionFirstHostDiskParallelLimit = "/admin/updateDecommissionFirstHostDiskParallelLimit" + AdminQueryDecommissionFirstHostDiskParallelLimit = "/admin/queryDecommissionFirstHostDiskParallelLimit" + AdminUpdateDecommissionFirstHostParallelLimit = "/admin/updateDecommissionFirstHostParallelLimit" + AdminQueryDecommissionFirstHostParallelLimit = "/admin/queryDecommissionFirstHostParallelLimit" + AdminQueryDecommissionFirstHostParallelInfo = "/admin/queryDecommissionFirstHostParallelInfo" + AdminUpdateDecommissionLimit = "/admin/updateDecommissionLimit" + AdminQueryDecommissionLimit = "/admin/queryDecommissionLimit" + AdminQueryDecommissionFailedDisk = "/admin/queryDecommissionFailedDisk" + AdminAbortDecommissionDisk = "/admin/abortDecommissionDisk" + AdminResetDataPartitionRestoreStatus = "/admin/resetDataPartitionRestoreStatus" + AdminGetOpLog = "/admin/getOpLog" + AdminGetRemoteCacheConfig = "/admin/getRemoteCacheConfig" // #nosec G101 AdminQueryDecommissionToken = "/admin/queryDecommissionToken" diff --git a/proto/model.go b/proto/model.go index 88eb9ef25..5b62fbae3 100644 --- a/proto/model.go +++ b/proto/model.go @@ -496,6 +496,13 @@ type DiscardDataPartitionInfos struct { DiscardDps []DataPartitionInfo } +type DecommissionStatusRecord struct { + Condition string + Status string + Time string + ErrMessage string +} + type DecommissionInfoStat struct { Key string RepairSourceDp []uint64 diff --git a/sdk/master/api_admin.go b/sdk/master/api_admin.go index ef1cc1e7e..9228b75d8 100644 --- a/sdk/master/api_admin.go +++ b/sdk/master/api_admin.go @@ -253,6 +253,14 @@ func (api *AdminAPI) AddMetaReplica(metaPartitionID uint64, nodeAddr string, cli return } +func (api *AdminAPI) QueryDataPartitionDecommissionStatusUpdateRecords(partitionId uint64) (records []*proto.DecommissionStatusRecord, err error) { + request := newRequest(get, proto.AdminQueryDataPartitionDecommissionStatusUpdateRecords).Header(api.h) + request.addParam("id", strconv.FormatUint(partitionId, 10)) + records = make([]*proto.DecommissionStatusRecord, 0) + err = api.mc.requestWith(&records, request) + return +} + func (api *AdminAPI) QueryDataPartitionDecommissionStatus(partitionId uint64) (info *proto.DecommissionDataPartitionInfo, err error) { request := newRequest(get, proto.AdminQueryDataPartitionDecommissionStatus).Header(api.h) request.addParam("id", strconv.FormatUint(partitionId, 10))