diff --git a/datanode/partition.go b/datanode/partition.go index e00e836cf..4f44d7385 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -1470,13 +1470,15 @@ func (dp *DataPartition) reload(s *SpaceManager) error { disk := dp.disk rootDir := dp.path log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir) - s.partitionMutex.Lock() - delete(s.partitions, dp.partitionID) - s.partitionMutex.Unlock() + s.DetachDataPartition(dp.partitionID) dp.Stop() dp.Disk().DetachDataPartition(dp) log.LogDebugf("data partition %v is detached", dp.partitionID) - _, err := LoadDataPartition(rootDir, disk) + dp2, err := LoadDataPartition(rootDir, disk) + if err != nil { + return err + } + s.AttachPartition(dp2) return err } diff --git a/datanode/server_handler.go b/datanode/server_handler.go index f803ab499..91f630f8a 100644 --- a/datanode/server_handler.go +++ b/datanode/server_handler.go @@ -767,3 +767,171 @@ func (s *DataNode) markDiskBroken(w http.ResponseWriter, r *http.Request) { } s.buildSuccessResp(w, "success") } + +func (s *DataNode) setDiskExtentReadLimitStatus(w http.ResponseWriter, r *http.Request) { + const ( + paramStatus = "status" + ) + if err := r.ParseForm(); err != nil { + err = fmt.Errorf("parse form fail: %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + status, err := strconv.ParseBool(r.FormValue(paramStatus)) + if err != nil { + err = fmt.Errorf("parse param %v fail: %v", paramStatus, err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + for _, disk := range s.space.disks { + disk.SetExtentRepairReadLimitStatus(status) + } + s.buildSuccessResp(w, "success") +} + +type DiskExtentReadLimitInfo struct { + DiskPath string `json:"diskPath"` + ExtentReadLimitStatus bool `json:"extentReadLimitStatus"` + Dp uint64 `json:"dp"` +} + +type DiskExtentReadLimitStatusResponse struct { + Infos []DiskExtentReadLimitInfo `json:"infos"` +} + +func (s *DataNode) queryDiskExtentReadLimitStatus(w http.ResponseWriter, r *http.Request) { + resp := &DiskExtentReadLimitStatusResponse{} + for _, disk := range s.space.disks { + status, dp := disk.QueryExtentRepairReadLimitStatus() + resp.Infos = append(resp.Infos, DiskExtentReadLimitInfo{DiskPath: disk.Path, ExtentReadLimitStatus: status, Dp: dp}) + } + s.buildSuccessResp(w, resp) +} + +func (s *DataNode) detachDataPartition(w http.ResponseWriter, r *http.Request) { + const ( + paramID = "id" + ) + if err := r.ParseForm(); err != nil { + err = fmt.Errorf("parse form fail: %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + partitionID, err := strconv.ParseUint(r.FormValue(paramID), 10, 64) + if err != nil { + err = fmt.Errorf("parse param %v fail: %v", paramID, err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + partition := s.space.Partition(partitionID) + if partition == nil { + s.buildFailureResp(w, http.StatusBadRequest, "partition not exist") + return + } + // store disk path and root of dp + disk := partition.disk + rootDir := partition.path + log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir) + + s.space.partitionMutex.Lock() + delete(s.space.partitions, partitionID) + s.space.partitionMutex.Unlock() + partition.Stop() + partition.Disk().DetachDataPartition(partition) + + log.LogDebugf("data partition %v is detached", partitionID) + s.buildSuccessResp(w, "success") +} + +func (s *DataNode) releaseDiskExtentReadLimitToken(w http.ResponseWriter, r *http.Request) { + const ( + paramDisk = "disk" + ) + if err := r.ParseForm(); err != nil { + err = fmt.Errorf("parse form fail: %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + diskPath := r.FormValue(paramDisk) + // store disk path and root of dp + disk, err := s.space.GetDisk(diskPath) + if err != nil { + log.LogErrorf("action[loadDataPartition] disk(%v) is not found err(%v).", diskPath, err) + s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("disk %v is not found", diskPath)) + return + } + disk.ReleaseReadExtentToken() + s.buildSuccessResp(w, "success") +} + +func (s *DataNode) loadDataPartition(w http.ResponseWriter, r *http.Request) { + const ( + paramID = "id" + paramDisk = "disk" + ) + if err := r.ParseForm(); err != nil { + err = fmt.Errorf("parse form fail: %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + partitionID, err := strconv.ParseUint(r.FormValue(paramID), 10, 64) + if err != nil { + err = fmt.Errorf("parse param %v fail: %v", paramID, err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + partition := s.space.Partition(partitionID) + if partition != nil { + s.buildFailureResp(w, http.StatusBadRequest, "partition is already loaded") + return + } + diskPath := r.FormValue(paramDisk) + // store disk path and root of dp + disk, err := s.space.GetDisk(diskPath) + if err != nil { + log.LogErrorf("action[loadDataPartition] disk(%v) is not found err(%v).", diskPath, err) + s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("disk %v is not found", diskPath)) + return + } + + fileInfoList, err := os.ReadDir(disk.Path) + if err != nil { + log.LogErrorf("action[loadDataPartition] read dir(%v) err(%v).", disk.Path, err) + s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf(" read dir(%v) err(%v)", disk.Path, err)) + return + } + rootDir := "" + for _, fileInfo := range fileInfoList { + filename := fileInfo.Name() + if !disk.isPartitionDir(filename) { + if disk.isExpiredPartitionDir(filename) { + } + continue + } + + if id, _, err := unmarshalPartitionName(filename); err != nil { + log.LogErrorf("action[RestorePartition] unmarshal partitionName(%v) from disk(%v) err(%v) ", + filename, disk.Path, err.Error()) + continue + } else { + if id == partitionID { + rootDir = filename + } + } + } + if rootDir == "" { + log.LogErrorf("action[loadDataPartition] dp root not found in dir(%v) .", disk.Path) + s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("dp root not found in dir(%v)", disk.Path)) + return + } + + log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir) + + dp, err := LoadDataPartition(path.Join(diskPath, rootDir), disk) + if err != nil { + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + } else { + s.space.AttachPartition(dp) + s.buildSuccessResp(w, "success") + } +} diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index 0383a569c..daa082f92 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -1325,7 +1325,8 @@ func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) { return } - log.LogInfof("action[handlePacketToAddDataPartitionRaftMember] %v, partition id %v", req.AddPeer, req.PartitionId) + log.LogInfof("action[handlePacketToAddDataPartitionRaftMember]req(%v) addPeer %v, partition id %v", + p.GetReqID(), req.AddPeer, req.PartitionId) p.AddMesgLog(string(reqData)) dp := s.space.Partition(req.PartitionId) @@ -1348,11 +1349,11 @@ func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) { isRaftLeader, err = s.forwardToRaftLeader(dp, p, false) if !isRaftLeader { if err != nil { - log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v forward to leader failed:%v", - dp.partitionID, p.GetReqID(), err) + log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v addPeer %v forward to leader failed:%v", + dp.partitionID, p.GetReqID(), req.AddPeer, err) } else { - log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v forward to leader", - dp.partitionID, p.GetReqID()) + log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v addPeer %v forward to leader", + dp.partitionID, p.GetReqID(), req.AddPeer) } return } @@ -1408,8 +1409,9 @@ func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) { p.GetReqID(), string(reqData), req.RemovePeer.Addr, dp.partitionID, dp.replicaNum, dp.config.Peers, dp.replicas) p.PartitionID = req.PartitionId - // do not return error to keep decommission progress go forward - // do not check replica existence on leader for autoRemove enable, follower may be contains redundant peers + // do not return error to keep master decommission progress go forward + // do not check replica existence on leader for autoRemove enable, follower may be contains redundant peers, make sure + // leader can send remove wal logs to follower if !dp.IsExistReplica(req.RemovePeer.Addr) && !req.Force && !req.AutoRemove { log.LogWarnf("action[handlePacketToRemoveDataPartitionRaftMember]dp %v receive MasterCommand: req %v "+ "RemoveRaftPeer(%v) force(%v) autoRemove(%v) has not exist", dp.partitionID, p.GetReqID(), req.RemovePeer, req.Force, req.AutoRemove) @@ -1456,7 +1458,8 @@ func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) { } } } - if !found && !req.AutoRemove { + + if !found && !req.AutoRemove && !req.Force { err = errors.NewErrorf("cannot found peer(%v) in dp(%v) peers", req.RemovePeer.Addr, dp.partitionID) log.LogWarnf("handlePacketToRemoveDataPartitionRaftMember:%v", err.Error()) return diff --git a/depends/tiglabs/raft/raft.go b/depends/tiglabs/raft/raft.go index 4bce9978d..c571208a7 100644 --- a/depends/tiglabs/raft/raft.go +++ b/depends/tiglabs/raft/raft.go @@ -211,6 +211,7 @@ func (s *raft) runApply() { } s.doStop() s.resetApply() + log.LogWarnf("raft(%v) quit runApply", s.raftFsm.id) }() loopCount := 0 @@ -271,6 +272,7 @@ func (s *raft) run() { s.stopSnapping() s.raftConfig.Storage.Close() close(s.done) + log.LogWarnf("raft(%v) quit run", s.raftFsm.id) }() s.prevHardSt.Term = s.raftFsm.term diff --git a/depends/tiglabs/raft/raft_fsm_candidate.go b/depends/tiglabs/raft/raft_fsm_candidate.go index 7ad4e4a13..05283d451 100644 --- a/depends/tiglabs/raft/raft_fsm_candidate.go +++ b/depends/tiglabs/raft/raft_fsm_candidate.go @@ -121,8 +121,8 @@ func (r *raftFsm) campaign(force bool, t CampaignType) { } li, lt := r.raftLog.lastIndexAndTerm() if logger.IsEnableDebug() { - logger.Debug("[raft->campaign][%v,%v logterm: %d, index: %d] sent "+ - "%v request to %v at term %d. raftFSM[%p]", msgType, r.id, r.config.ReplicateAddr, lt, li, id, r.term, r) + logger.Debug("[raft->campaign][%v,raft %v, term: %d] sent "+ + "index %v request to %v at term %d. raftFSM[%p]", msgType, r.id, lt, li, id, r.term, r) } m := proto.GetMessage() diff --git a/depends/tiglabs/raft/server.go b/depends/tiglabs/raft/server.go index cdda9eb7c..49e40827c 100644 --- a/depends/tiglabs/raft/server.go +++ b/depends/tiglabs/raft/server.go @@ -234,6 +234,8 @@ func (rs *RaftServer) Status(id uint64) (status *Status) { if ok { status = raft.status() + } else { + logger.Warn("raftServer cannot found, id:%d", id) } if status == nil { status = &Status{ diff --git a/master/cluster.go b/master/cluster.go index 352cb93a1..b8079d1f6 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -2473,6 +2473,9 @@ func (c *Cluster) addDataReplica(dp *DataPartition, addr string, ignoreDecommiss // update datanode size with to replica size func (c *Cluster) updateDataNodeSize(addr string, dp *DataPartition) error { + if len(dp.Replicas) == 0 { + return errors.NewErrorf("dp %v has empty replica", dp.decommissionInfo()) + } leaderSize := dp.Replicas[0].Used dataNode, err := c.dataNode(addr) if err != nil { @@ -2576,7 +2579,17 @@ func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Pe _, err = c.buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp, addPeer, host) if err == nil { break + } else { + // if send to leader raise err, it may send to follower ,then follower forward + // this request to leader, return nil. so when leader encounter en error, should + // return err + if leaderAddr != "" && leaderAddr == host { + dp.Hosts = oldHosts + dp.Peers = oldPeers + return err + } } + if index < len(candidateAddrs)-1 { time.Sleep(retrySendSyncTaskInternal) } diff --git a/master/data_node.go b/master/data_node.go index 76c362702..b8560ffc3 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -16,6 +16,7 @@ package master import ( "fmt" + "github.com/cubefs/cubefs/util/auditlog" "sync" "sync/atomic" "time" @@ -99,10 +100,12 @@ func (dataNode *DataNode) SetIoUtils(used map[string]float64) { func (dataNode *DataNode) checkLiveness() { dataNode.Lock() defer dataNode.Unlock() - log.LogInfof("action[checkLiveness] datanode[%v] report time[%v],since report time[%v], need gap [%v]", - dataNode.Addr, dataNode.ReportTime, time.Since(dataNode.ReportTime), time.Second*time.Duration(defaultNodeTimeOutSec)) if time.Since(dataNode.ReportTime) > time.Second*time.Duration(defaultNodeTimeOutSec) { dataNode.isActive = false + msg := fmt.Sprintf("datanode[%v] report time[%v],since report time[%v], need gap [%v]", + dataNode.Addr, dataNode.ReportTime, time.Since(dataNode.ReportTime), time.Second*time.Duration(defaultNodeTimeOutSec)) + log.LogWarnf("action[checkLiveness] %v", msg) + auditlog.LogMasterOp("DataNodeLive", msg, nil) } } diff --git a/master/data_partition.go b/master/data_partition.go index d84c3ed75..b38fef26b 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -616,9 +616,12 @@ func (partition *DataPartition) getLiveReplicasFromHosts(timeOutSec int64) (repl if replica.isLive(partition.PartitionID, timeOutSec) { replicas = append(replicas, replica) } else { + msg := fmt.Sprintf("dp %v replica addr %v is unavailable, datanode active %v replica status %v and is active %v", + partition.PartitionID, replica.Addr, replica.dataNode.isActive, replica.Status, replica.isActive(timeOutSec)) replica.Status = proto.Unavailable log.LogWarnf("action[getLiveReplicasFromHosts] vol %v dp %v replica %v is unavailable", partition.VolName, partition.PartitionID, replica.Addr) + auditlog.LogMasterOp("DataPartitionReplicaStatus", msg, nil) } } @@ -2044,7 +2047,8 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) { } // remove raft member err = partition.createTaskToRemoveRaftMember(c, peer, force, true) - auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v force %v", partition.PartitionID, peer, force) + auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v force %v:to replica %v: LocalPeers%v", + partition.decommissionInfo(), peer, force, replica.Addr, replica.LocalPeers) log.LogDebugf("action[checkReplicaMeta]%v, err %v", auditMsg, err) auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err) if err != nil { @@ -2061,8 +2065,25 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) { redundantPeers := findPeersToDeleteByConfig(partition.Peers, replica.LocalPeers) for _, peer := range redundantPeers { err = c.removeHostMember(partition, peer) - auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v for master", partition.PartitionID, peer) - log.LogDebugf("action[checkReplicaMeta]%v: err %v", auditMsg, err) + auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v for master,base on replica %v,localPeers(%v) ", + partition.decommissionInfo(), peer, replica.Addr, replica.LocalPeers) + auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err) + if err != nil { + return + } + // redundant peers on master may exist on dataNode, and the redundant replica will be + // added into partition.Replicas again by hear beat. + var dataNode *DataNode + dataNode, err = c.dataNode(peer.Addr) + auditMsg = fmt.Sprintf("dp(%v) cannot found datanode for replica %v ,base on replica %v,localPeers(%v) ", + partition.decommissionInfo(), peer.Addr, replica.Addr, replica.LocalPeers) + auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err) + if err != nil { + return + } + err = c.deleteDataReplica(partition, dataNode) + auditMsg = fmt.Sprintf("dp(%v) remove redundant replica on %v for master,base on replica %v,localPeers(%v) ", + partition.decommissionInfo(), peer.Addr, replica.Addr, replica.LocalPeers) auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err) if err != nil { return @@ -2133,9 +2154,9 @@ func (partition *DataPartition) lostLeader(c *Cluster) bool { } func (partition *DataPartition) decommissionInfo() string { - return fmt.Sprintf("vol(%v)_dp(%v)_src(%v)_dst(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+ + 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)", - partition.VolName, partition.PartitionID, partition.DecommissionSrcAddr, partition.DecommissionDstAddr, + 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), @@ -2225,8 +2246,14 @@ func (partition *DataPartition) setRestoreReplicaStop() bool { } func (partition *DataPartition) tryRestoreReplicaMeta(c *Cluster, migrateType uint32) error { - // AutoAddReplica do not need to check meta for replica again + // AutoAddReplica do not need to check meta for replica again, only have to check + // dp is performing decommission if migrateType == AutoAddReplica { + if partition.isPerformingDecommission(c) { + log.LogDebugf("action[checkReplicaMeta]dp(%v) is performing decommission, skip it", + partition.PartitionID) + return proto.ErrPerformingDecommission + } return nil } // diff --git a/master/disk_manager.go b/master/disk_manager.go index 39c41533f..430d1c005 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -191,8 +191,9 @@ func (c *Cluster) checkDiskRecoveryProgress() { } else { partition.DecommissionErrorMessage = "" partition.SetDecommissionStatus(DecommissionSuccess) // can be readonly or readwrite - Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered success", - c.Name, partitionID, partition.DecommissionDstAddr)) + 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())) } partition.RLock() err = c.syncUpdateDataPartition(partition) diff --git a/master/topology.go b/master/topology.go index d6e26de8a..ffcc5392e 100644 --- a/master/topology.go +++ b/master/topology.go @@ -17,6 +17,7 @@ package master import ( "container/list" "fmt" + "github.com/cubefs/cubefs/util/auditlog" "sort" "strings" "sync" @@ -2191,6 +2192,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { if dp.IsDecommissionSuccess() { l.Remove(dp) dp.ReleaseDecommissionToken(c) + msg := fmt.Sprintf("dp %v decommission success, cost %v", + dp.decommissionInfo(), time.Since(dp.RecoverStartTime)) dp.ResetDecommissionStatus() dp.setRestoreReplicaStop() err := c.syncUpdateDataPartition(dp) @@ -2201,6 +2204,7 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for success", dp.PartitionID) } + auditlog.LogMasterOp("TraverseDataPartition", msg, err) } else if dp.IsDecommissionFailed() { if !dp.tryRollback(c) { log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for fail", @@ -2212,6 +2216,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { // rollback fail/success need release token dp.ReleaseDecommissionToken(c) c.syncUpdateDataPartition(dp) + msg := fmt.Sprintf("dp %v decommission failed", dp.decommissionInfo()) + auditlog.LogMasterOp("TraverseDataPartition", msg, nil) } else if dp.IsDecommissionPaused() { log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for paused ", dp.PartitionID)