diff --git a/cli/cmd/datapartition.go b/cli/cmd/datapartition.go index 43bc6c995..3834d8cf8 100644 --- a/cli/cmd/datapartition.go +++ b/cli/cmd/datapartition.go @@ -278,6 +278,7 @@ The "reset" command will be released in next version`, func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command { var raftForceDel bool + var decommissionType string var clientIDKey string cmd := &cobra.Command{ Use: CliOpDecommission + " [ADDRESS] [DATA PARTITION ID]", @@ -296,7 +297,7 @@ func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command if err != nil { return } - if err := client.AdminAPI().DecommissionDataPartition(partitionID, address, raftForceDel, clientIDKey); err != nil { + if err := client.AdminAPI().DecommissionDataPartition(partitionID, address, raftForceDel, clientIDKey, decommissionType); err != nil { stdout(fmt.Sprintf("failed:err(%v)\n", err.Error())) return } @@ -310,6 +311,7 @@ func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command }, } cmd.Flags().BoolVarP(&raftForceDel, "raftForceDel", "r", false, "true for raftForceDel") + cmd.Flags().StringVar(&decommissionType, "decommissionType", "1", "decommission type") cmd.Flags().StringVar(&clientIDKey, CliFlagClientIDKey, client.ClientIDKey(), CliUsageClientIDKey) return cmd } diff --git a/datanode/partition.go b/datanode/partition.go index 859c7bd7d..f74336cf7 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -690,7 +690,6 @@ func (dp *DataPartition) ForceLoadHeader() { func (dp *DataPartition) RemoveAll(decommissionType uint32, force, isSpecialReplica bool) (err error) { dp.persistMetaMutex.Lock() defer dp.persistMetaMutex.Unlock() - if force && isSpecialReplica && decommissionType == proto.AutoDecommission { originalPath := dp.Path() parent := path.Dir(originalPath) @@ -754,6 +753,7 @@ func (dp *DataPartition) PersistMetadata() (err error) { if metaData, err = json.Marshal(md); err != nil { return } + // persist meta can be failed with io error if _, err = metadataFile.Write(metaData); err != nil { return } @@ -1481,8 +1481,8 @@ func (dp *DataPartition) isDecommissionRecovering() bool { func (dp *DataPartition) incDiskErrCnt() { diskErrCnt := atomic.AddUint64(&dp.diskErrCnt, 1) - dp.PersistMetadata() - log.LogWarnf("[incDiskErrCnt]: dp(%v) disk err count:%v", dp.partitionID, diskErrCnt) + err := dp.PersistMetadata() + log.LogWarnf("[incDiskErrCnt]: dp(%v) disk err count:%v, err %v", dp.partitionID, diskErrCnt, err) } func (dp *DataPartition) getDiskErrCnt() uint64 { diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index fa445f822..cde534620 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -21,6 +21,8 @@ import ( "fmt" "hash/crc32" "net" + "os" + "path" "strconv" "strings" "sync" @@ -191,6 +193,8 @@ func (s *DataNode) OperatePacket(p *repl.Packet, c net.Conn) (err error) { s.handlePacketToStopDataPartitionRepair(p) case proto.OpRecoverDataReplicaMeta: s.handlePacketToRecoverDataReplicaMeta(p) + case proto.OpRecoverBackupDataReplica: + s.handlePacketToRecoverBackupDataReplica(p) default: p.PackErrorBody(repl.ErrorUnknownOp.Error(), repl.ErrorUnknownOp.Error()+strconv.Itoa(int(p.Opcode))) } @@ -1770,3 +1774,84 @@ func (s *DataNode) handleBatchUnlockNormalExtent(p *repl.Packet, connect net.Con log.LogInfof("action[handleBatchUnlockNormalExtent] success len: %v", len(exts)) } + +func (s *DataNode) handlePacketToRecoverBackupDataReplica(p *repl.Packet) { + task := &proto.AdminTask{} + err := json.Unmarshal(p.Data, task) + defer func() { + if err != nil { + p.PackErrorBody(ActionRecoverDataReplicaMeta, err.Error()) + } else { + p.PacketOkReply() + } + }() + if err != nil { + return + } + request := &proto.RecoverBackupDataReplicaRequest{} + if task.OpCode != proto.OpRecoverBackupDataReplica { + err = fmt.Errorf("action[handlePacketToRecoverBackupDataReplica] illegal opcode ") + log.LogWarnf("action[handlePacketToRecoverBackupDataReplica] illegal opcode ") + return + } + + bytes, _ := json.Marshal(task.Request) + p.AddMesgLog(string(bytes)) + err = json.Unmarshal(bytes, request) + if err != nil { + return + } + log.LogDebugf("action[handlePacketToRecoverBackupDataReplica] try recover %v", request.PartitionId) + disk, err := s.space.GetDisk(request.Disk) + if err != nil { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] disk(%v) is not found err(%v).", request.Disk, err) + return + } + + fileInfoList, err := os.ReadDir(disk.Path) + if err != nil { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] read dir(%v) err(%v).", disk.Path, err) + return + } + rootDir := "" + for _, fileInfo := range fileInfoList { + filename := fileInfo.Name() + + if !disk.isBackupPartitionDir(filename) { + continue + } + + if id, err := unmarshalBackupPartitionDirName(filename); err != nil { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] unmarshal partitionName(%v) from disk(%v) err(%v) ", + filename, disk.Path, err.Error()) + continue + } else { + if id == request.PartitionId { + rootDir = filename + } + } + } + if rootDir == "" { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] dp root not found in dir(%v) .", disk.Path) + return + } + + // rename root dir back to normal + newPath := strings.Replace(rootDir, BackupPartitionPrefix, "", 1) + err = os.Rename(path.Join(disk.Path, rootDir), path.Join(disk.Path, newPath)) + if err != nil { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] rename disk %v rootDir %v to %v failed %v.", + disk.Path, rootDir, newPath, err) + return + } + + dp, err := LoadDataPartition(path.Join(request.Disk, newPath), disk) + if err != nil { + log.LogErrorf("action[handlePacketToRecoverBackupDataReplica] load disk %v rootDir %v failed err %v.", + disk.Path, rootDir, err) + } else { + s.space.AttachPartition(dp) + log.LogInfof("action[handlePacketToRecoverBackupDataReplica] load disk %v rootDir %v success .", + disk.Path, rootDir) + } +} diff --git a/master/api_service.go b/master/api_service.go index 389b1d276..5a7125e3d 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -1627,6 +1627,21 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) { return } + retry := 0 + for { + if !dp.setRestoreReplicaForbidden() { + retry++ + if retry > defaultDecommissionRetryLimit { + } else { + err = errors.NewErrorf("set RestoreReplicaMetaForbidden failed") + sendErrReply(w, r, newErrHTTPReply(err)) + return + } + time.Sleep(1 * time.Second) + } + break + } + if err = m.cluster.addDataReplica(dp, addr, false); err != nil { sendErrReply(w, r, newErrHTTPReply(err)) return @@ -1636,7 +1651,6 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) { dp.DecommissionType = ManualAddReplica dp.RecoverStartTime = time.Now() dp.SetDecommissionStatus(DecommissionRunning) - dp.setRestoreReplicaForbidden() dp.Status = proto.ReadOnly dp.isRecover = true @@ -1895,6 +1909,15 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) return } + decommissionType, err := parseUintParam(r, DecommissionType) + if err != nil { + sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) + return + } + // default is ManualDecommission + if decommissionType == 0 { + decommissionType = int(ManualDecommission) + } node, err := c.dataNode(addr) if err != nil { rstMsg = fmt.Sprintf(" dataPartitionID :%v not find datanode for addr %v", @@ -1902,7 +1925,7 @@ 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, raftForce, ManualDecommission) + err = m.cluster.markDecommissionDataPartition(dp, node, raftForce, uint32(decommissionType)) if err != nil { rstMsg = err.Error() sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: rstMsg}) @@ -7053,3 +7076,87 @@ func (m *Server) getAllMetaNodes(w http.ResponseWriter, r *http.Request) { metaNodes := m.cluster.allMetaNodes() sendOkReply(w, r, newSuccessHTTPReply(metaNodes)) } + +func (m *Server) recoverDiskErrorReplica(w http.ResponseWriter, r *http.Request) { + var ( + msg string + addr string + dp *DataPartition + dataNode *DataNode + partitionID uint64 + err error + ) + metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminRecoverDiskErrorReplica)) + defer func() { + doStatAndMetric(proto.AdminRecoverDiskErrorReplica, metric, err, nil) + }() + + if partitionID, addr, err = parseRequestToAddDataReplica(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 + } + + dataNode, err = m.cluster.dataNode(addr) + if err != nil { + sendErrReply(w, r, newErrHTTPReply(err)) + return + } + + if !proto.IsNormalDp(dp.PartitionType) { + err = fmt.Errorf("action[recoverDiskErrorReplica] [%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) + sendErrReply(w, r, newErrHTTPReply(err)) + } + + retry := 0 + for { + if !dp.setRestoreReplicaForbidden() { + retry++ + if retry > defaultDecommissionRetryLimit { + } else { + err = errors.NewErrorf("set RestoreReplicaMetaForbidden failed") + sendErrReply(w, r, newErrHTTPReply(err)) + return + } + time.Sleep(1 * time.Second) + } + break + } + defer dp.setRestoreReplicaStop() + // 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) + 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) + 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) + 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) + sendErrReply(w, r, newErrHTTPReply(err)) + return + } + msg = fmt.Sprintf("action[recoverDiskErrorReplica] dp(%v) recover replica [%v_%v] successfully", dp.decommissionInfo(), backupInfo.Addr, backupInfo.Disk) + log.LogInfof("%v", msg) + sendOkReply(w, r, newSuccessHTTPReply(msg)) +} diff --git a/master/cluster.go b/master/cluster.go index d5ca7c7cd..922053af3 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -2787,7 +2787,8 @@ func (c *Cluster) deleteDataReplica(dp *DataPartition, dataNode *DataNode, raftF // in case dataNode is unreachable,update meta first. dp.removeReplicaByAddr(dataNode.Addr) dp.checkAndRemoveMissReplica(dataNode.Addr) - + log.LogDebugf("action[deleteDataReplica] vol[%v],data partition[%v] remove replica[%v] force(%v)", + dp.VolName, dp.decommissionInfo(), dataNode.Addr, raftForceDel) if err = dp.update("deleteDataReplica", dp.VolName, dp.Peers, dp.Hosts, c); err != nil { dp.Unlock() return @@ -5019,3 +5020,18 @@ func (c *Cluster) RetryDecommissionDisk(addr string, diskPath string) bool { } return false } + +func (c *Cluster) syncRecoverBackupDataPartitionReplica(host, disk string, dp *DataPartition) (err error) { + log.LogInfof("action[syncRecoverBackupDataPartitionReplica] dp [%v] to recover replica on %v_%v", dp.PartitionID, host, disk) + var dataNode *DataNode + dataNode, err = c.dataNode(host) + if err != nil { + return + } + var task *proto.AdminTask + task = dp.createTaskToRecoverBackupDataPartitionReplica(host, disk) + if _, err = dataNode.TaskManager.syncSendAdminTask(task); err != nil { + return + } + return +} diff --git a/master/data_node.go b/master/data_node.go index a04f7394b..e2a0f7307 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -24,6 +24,7 @@ import ( "github.com/cubefs/cubefs/util" "github.com/cubefs/cubefs/util/atomicutil" "github.com/cubefs/cubefs/util/auditlog" + "github.com/cubefs/cubefs/util/errors" "github.com/cubefs/cubefs/util/log" ) @@ -571,3 +572,15 @@ func (dataNode *DataNode) getBackupDataPartitionIDs() (ids []uint64) { } return ids } + +func (dataNode *DataNode) getBackupDataPartitionInfo(id uint64) (proto.BackupDataPartitionInfo, error) { + dataNode.RLock() + dataNode.RUnlock() + for _, info := range dataNode.BackupDataPartitions { + if info.PartitionID == id { + return info, nil + } + } + return proto.BackupDataPartitionInfo{}, errors.NewErrorf("cannot find backup info "+ + "for dp (%v) on datanode (%v)", id, dataNode.Addr) +} diff --git a/master/data_partition.go b/master/data_partition.go index b9954ea10..f123a82a6 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -2289,3 +2289,11 @@ func (partition *DataPartition) getSpecifyStatusReplicaNum(status int8) uint8 { } return count } + +func (partition *DataPartition) createTaskToRecoverBackupDataPartitionReplica(addr, disk string) (task *proto.AdminTask, +) { + task = proto.NewAdminTask(proto.OpRecoverBackupDataReplica, addr, newRecoverBackupDataPartitionReplicaRequest( + partition.PartitionID, disk)) + partition.resetTaskID(task) + return +} diff --git a/master/http_server.go b/master/http_server.go index 2857f0e32..83bd9fd8f 100644 --- a/master/http_server.go +++ b/master/http_server.go @@ -565,6 +565,9 @@ func (m *Server) registerAPIRoutes(router *mux.Router) { router.NewRoute().Methods(http.MethodGet). Path(proto.AdminRecoverReplicaMeta). HandlerFunc(m.recoverReplicaMeta) + router.NewRoute().Methods(http.MethodGet). + Path(proto.AdminRecoverDiskErrorReplica). + HandlerFunc(m.recoverDiskErrorReplica) // meta node management APIs router.NewRoute().Methods(http.MethodGet, http.MethodPost). diff --git a/master/operate_util.go b/master/operate_util.go index a4787cba8..45a09e6de 100644 --- a/master/operate_util.go +++ b/master/operate_util.go @@ -262,3 +262,11 @@ func matchKey(serverKey, clientKey string) bool { cipherStr := h.Sum(nil) return strings.EqualFold(clientKey, hex.EncodeToString(cipherStr)) } + +func newRecoverBackupDataPartitionReplicaRequest(ID uint64, disk string) (req *proto.RecoverBackupDataReplicaRequest) { + req = &proto.RecoverBackupDataReplicaRequest{ + PartitionId: ID, + Disk: disk, + } + return +} diff --git a/proto/admin_proto.go b/proto/admin_proto.go index acd340455..ba4519f9f 100644 --- a/proto/admin_proto.go +++ b/proto/admin_proto.go @@ -1300,3 +1300,8 @@ type BackupDataPartitionInfo struct { Disk string PartitionID uint64 } + +type RecoverBackupDataReplicaRequest struct { + PartitionId uint64 + Disk string +} diff --git a/proto/packet.go b/proto/packet.go index f4b6622ae..a8c0d74ed 100644 --- a/proto/packet.go +++ b/proto/packet.go @@ -150,6 +150,7 @@ const ( OpQos uint8 = 0x6A OpStopDataPartitionRepair uint8 = 0x6B OpRecoverDataReplicaMeta uint8 = 0x6C + OpRecoverBackupDataReplica uint8 = 0x6D // Operations: MultipartInfo OpCreateMultipart uint8 = 0x70 diff --git a/repl/packet.go b/repl/packet.go index bfc2dc683..180aee3ee 100644 --- a/repl/packet.go +++ b/repl/packet.go @@ -494,7 +494,8 @@ func (p *Packet) IsMasterCommand() bool { proto.OpDecommissionDataPartition, proto.OpAddDataPartitionRaftMember, proto.OpRemoveDataPartitionRaftMember, - proto.OpDataPartitionTryToLeader: + proto.OpDataPartitionTryToLeader, + proto.OpRecoverBackupDataReplica: return true default: return false diff --git a/sdk/master/api_admin.go b/sdk/master/api_admin.go index 0cb11bc47..4d55abab6 100644 --- a/sdk/master/api_admin.go +++ b/sdk/master/api_admin.go @@ -170,12 +170,13 @@ func (api *AdminAPI) CreateDataPartition(volName string, count int, clientIDKey )) } -func (api *AdminAPI) DecommissionDataPartition(dataPartitionID uint64, nodeAddr string, raftForce bool, clientIDKey string) (err error) { +func (api *AdminAPI) DecommissionDataPartition(dataPartitionID uint64, nodeAddr string, raftForce bool, clientIDKey, decommissionType string) (err error) { request := newRequest(get, proto.AdminDecommissionDataPartition).Header(api.h) request.addParam("id", strconv.FormatUint(dataPartitionID, 10)) request.addParam("addr", nodeAddr) request.addParam("raftForceDel", strconv.FormatBool(raftForce)) request.addParam("clientIDKey", clientIDKey) + request.addParam("decommissionType", decommissionType) _, err = api.mc.serveRequest(request) return }