mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
style(datanode): format code of datanode
Signed-off-by: slasher <shenjie1@oppo.com>
This commit is contained in:
parent
5e93d2cea8
commit
cf415f1acf
@ -158,7 +158,7 @@ func (dp *DataPartition) buildDataPartitionRepairTask(repairTasks []*DataPartiti
|
||||
}
|
||||
|
||||
func (dp *DataPartition) getLocalExtentInfo(extentType uint8, tinyExtents []uint64) (extents []*storage.ExtentInfo, leaderTinyDeleteRecordFileSize int64, err error) {
|
||||
localExtents := make([]*storage.ExtentInfo, 0)
|
||||
var localExtents []*storage.ExtentInfo
|
||||
|
||||
if extentType == proto.NormalExtentType {
|
||||
localExtents, leaderTinyDeleteRecordFileSize, err = dp.extentStore.GetAllWatermarks(storage.NormalExtentFilter())
|
||||
@ -175,8 +175,7 @@ func (dp *DataPartition) getLocalExtentInfo(extentType uint8, tinyExtents []uint
|
||||
}
|
||||
extents = make([]*storage.ExtentInfo, 0, len(localExtents))
|
||||
for _, et := range localExtents {
|
||||
newEt := storage.ExtentInfo{}
|
||||
newEt = *et
|
||||
newEt := *et
|
||||
extents = append(extents, &newEt)
|
||||
}
|
||||
return
|
||||
@ -259,7 +258,6 @@ func (dp *DataPartition) moveToBrokenTinyExtentC(extentType uint8, extents []uin
|
||||
if extentType == proto.TinyExtentType {
|
||||
dp.extentStore.SendAllToBrokenTinyExtentC(extents)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (dp *DataPartition) sendAllTinyExtentsToC(extentType uint8, availableTinyExtents, brokenTinyExtents []uint64) {
|
||||
@ -341,7 +339,7 @@ func (dp *DataPartition) buildExtentCreationTasks(repairTasks []*DataPartitionRe
|
||||
if repairTask == nil {
|
||||
continue
|
||||
}
|
||||
if _, ok := repairTask.extents[extentID]; !ok && extentInfo.IsDeleted == false {
|
||||
if _, ok := repairTask.extents[extentID]; !ok && !extentInfo.IsDeleted {
|
||||
if storage.IsTinyExtent(extentID) {
|
||||
continue
|
||||
}
|
||||
|
||||
@ -39,10 +39,10 @@ import (
|
||||
|
||||
var (
|
||||
// RegexpDataPartitionDir validates the directory name of a data partition.
|
||||
RegexpDataPartitionDir, _ = regexp.Compile("^datapartition_(\\d)+_(\\d)+$")
|
||||
RegexpCachePartitionDir, _ = regexp.Compile("^cachepartition_(\\d)+_(\\d)+$")
|
||||
RegexpPreLoadPartitionDir, _ = regexp.Compile("^preloadpartition_(\\d)+_(\\d)+$")
|
||||
RegexpExpiredDataPartitionDir, _ = regexp.Compile("^expired_datapartition_(\\d)+_(\\d)+$")
|
||||
RegexpDataPartitionDir, _ = regexp.Compile(`^datapartition_(\d)+_(\d)+$`)
|
||||
RegexpCachePartitionDir, _ = regexp.Compile(`^cachepartition_(\d)+_(\d)+$`)
|
||||
RegexpPreLoadPartitionDir, _ = regexp.Compile(`^preloadpartition_(\d)+_(\d)+$`)
|
||||
RegexpExpiredDataPartitionDir, _ = regexp.Compile(`^expired_datapartition_(\d)+_(\d)+$`)
|
||||
)
|
||||
|
||||
const (
|
||||
@ -336,7 +336,6 @@ func (d *Disk) triggerDiskError(err error) {
|
||||
d.ForceExitRaftStore()
|
||||
d.Status = proto.Unavailable
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (d *Disk) updateSpaceInfo() (err error) {
|
||||
@ -535,13 +534,11 @@ func (d *Disk) RestorePartition(visitor PartitionVisitor) (err error) {
|
||||
go func(toDeleteExpiredPartitions []string) {
|
||||
ticker := time.NewTicker(ExpiredPartitionExistTime)
|
||||
log.LogInfof("action[RestorePartition] delete expiredPartitions automatically start, toDeleteExpiredPartitions %v", toDeleteExpiredPartitions)
|
||||
select {
|
||||
case <-ticker.C:
|
||||
d.deleteExpiredPartitions(toDeleteExpiredPartitionNames)
|
||||
ticker.Stop()
|
||||
log.LogInfof("action[RestorePartition] delete expiredPartitions automatically finish")
|
||||
return
|
||||
}
|
||||
|
||||
<-ticker.C
|
||||
d.deleteExpiredPartitions(toDeleteExpiredPartitionNames)
|
||||
ticker.Stop()
|
||||
log.LogInfof("action[RestorePartition] delete expiredPartitions automatically finish")
|
||||
}(notDeletedExpiredPartitionNames)
|
||||
}
|
||||
}
|
||||
|
||||
@ -18,7 +18,7 @@ const (
|
||||
MaxRepairErrCnt = 1000
|
||||
)
|
||||
|
||||
var nodeInfoStopC = make(chan struct{}, 0)
|
||||
var nodeInfoStopC = make(chan struct{})
|
||||
|
||||
func (m *DataNode) startUpdateNodeInfo() {
|
||||
ticker := time.NewTicker(UpdateNodeInfoTicket)
|
||||
|
||||
@ -24,11 +24,9 @@ import (
|
||||
"os"
|
||||
"path"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
raftProto "github.com/cubefs/cubefs/depends/tiglabs/raft/proto"
|
||||
@ -72,20 +70,6 @@ type DataPartitionMetadata struct {
|
||||
StopRecover bool
|
||||
}
|
||||
|
||||
type sortedPeers []proto.Peer
|
||||
|
||||
func (sp sortedPeers) Len() int {
|
||||
return len(sp)
|
||||
}
|
||||
|
||||
func (sp sortedPeers) Less(i, j int) bool {
|
||||
return sp[i].ID < sp[j].ID
|
||||
}
|
||||
|
||||
func (sp sortedPeers) Swap(i, j int) {
|
||||
sp[i], sp[j] = sp[j], sp[i]
|
||||
}
|
||||
|
||||
func (md *DataPartitionMetadata) Validate() (err error) {
|
||||
md.VolumeID = strings.TrimSpace(md.VolumeID)
|
||||
if len(md.VolumeID) == 0 || md.PartitionID == 0 || md.PartitionSize == 0 {
|
||||
@ -318,8 +302,8 @@ func newDataPartition(dpCfg *dataPartitionCfg, disk *Disk, isCreate bool) (dp *D
|
||||
partitionSize: dpCfg.PartitionSize,
|
||||
partitionType: dpCfg.PartitionType,
|
||||
replicas: make([]string, 0),
|
||||
stopC: make(chan bool, 0),
|
||||
stopRaftC: make(chan uint64, 0),
|
||||
stopC: make(chan bool),
|
||||
stopRaftC: make(chan uint64),
|
||||
storeC: make(chan uint64, 128),
|
||||
snapshot: make([]*proto.File, 0),
|
||||
partitionStatus: proto.ReadWrite,
|
||||
@ -353,9 +337,7 @@ func (dp *DataPartition) replicasInit() {
|
||||
if dp.config.Hosts == nil {
|
||||
return
|
||||
}
|
||||
for _, host := range dp.config.Hosts {
|
||||
replicas = append(replicas, host)
|
||||
}
|
||||
replicas = append(replicas, dp.config.Hosts...)
|
||||
dp.replicasLock.Lock()
|
||||
dp.replicas = replicas
|
||||
dp.replicasLock.Unlock()
|
||||
@ -415,8 +397,7 @@ func (dp *DataPartition) getReplicaCopy() []string {
|
||||
dp.replicasLock.RLock()
|
||||
defer dp.replicasLock.RUnlock()
|
||||
|
||||
var tmpCopy []string
|
||||
tmpCopy = make([]string, len(dp.replicas))
|
||||
tmpCopy := make([]string, len(dp.replicas))
|
||||
copy(tmpCopy, dp.replicas)
|
||||
|
||||
return tmpCopy
|
||||
@ -434,24 +415,6 @@ func (dp *DataPartition) getReplicaLen() int {
|
||||
return len(dp.replicas)
|
||||
}
|
||||
|
||||
func (dp *DataPartition) needDeleteReplica(addr string) bool {
|
||||
if dp.IsExsitReplica(addr) {
|
||||
return true
|
||||
}
|
||||
|
||||
if dp.config == nil {
|
||||
return false
|
||||
}
|
||||
|
||||
for _, h := range dp.config.Hosts {
|
||||
if addr == h {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
func (dp *DataPartition) IsExsitReplica(addr string) bool {
|
||||
dp.replicasLock.RLock()
|
||||
defer dp.replicasLock.RUnlock()
|
||||
@ -497,7 +460,6 @@ func (dp *DataPartition) Stop() {
|
||||
dp.extentStore.Close()
|
||||
_ = dp.storeAppliedID(atomic.LoadUint64(&dp.appliedID))
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
// Disk returns the disk instance.
|
||||
@ -552,9 +514,6 @@ func (dp *DataPartition) PersistMetadata() (err error) {
|
||||
os.Remove(fileName)
|
||||
}()
|
||||
|
||||
// sp := sortedPeers(dp.config.Peers)
|
||||
// sort.Sort(sp)
|
||||
|
||||
md := &DataPartitionMetadata{
|
||||
VolumeID: dp.config.VolName,
|
||||
PartitionID: dp.config.PartitionID,
|
||||
@ -637,37 +596,6 @@ func (dp *DataPartition) statusUpdate() {
|
||||
dp.partitionStatus = int(math.Min(float64(status), float64(dp.disk.Status)))
|
||||
}
|
||||
|
||||
func parseFileName(filename string) (extentID uint64, isExtent bool) {
|
||||
if isExtent = storage.RegexpExtentFile.MatchString(filename); !isExtent {
|
||||
return
|
||||
}
|
||||
var err error
|
||||
if extentID, err = strconv.ParseUint(filename, 10, 64); err != nil {
|
||||
isExtent = false
|
||||
return
|
||||
}
|
||||
isExtent = true
|
||||
return
|
||||
}
|
||||
|
||||
func (dp *DataPartition) actualSize(path string, finfo os.FileInfo) (size int64) {
|
||||
name := finfo.Name()
|
||||
extentID, isExtent := parseFileName(name)
|
||||
if !isExtent {
|
||||
return 0
|
||||
}
|
||||
if storage.IsTinyExtent(extentID) {
|
||||
stat := new(syscall.Stat_t)
|
||||
err := syscall.Stat(fmt.Sprintf("%v/%v", path, finfo.Name()), stat)
|
||||
if err != nil {
|
||||
return finfo.Size()
|
||||
}
|
||||
return stat.Blocks * DiskSectorSize
|
||||
}
|
||||
|
||||
return finfo.Size()
|
||||
}
|
||||
|
||||
func (dp *DataPartition) computeUsage() {
|
||||
if time.Now().Unix()-dp.intervalToUpdatePartitionSize < IntervalToUpdatePartitionSize {
|
||||
return
|
||||
@ -755,19 +683,15 @@ func (dp *DataPartition) updateReplicas(isForce bool) (err error) {
|
||||
|
||||
// Compare the fetched replica with the local one.
|
||||
func (dp *DataPartition) compareReplicas(v1, v2 []string) (equals bool) {
|
||||
equals = true
|
||||
if len(v1) == len(v2) {
|
||||
for i := 0; i < len(v1); i++ {
|
||||
if v1[i] != v2[i] {
|
||||
equals = false
|
||||
return
|
||||
return false
|
||||
}
|
||||
}
|
||||
equals = true
|
||||
return
|
||||
return true
|
||||
}
|
||||
equals = false
|
||||
return
|
||||
return false
|
||||
}
|
||||
|
||||
// Fetch the replica information from the master.
|
||||
@ -787,9 +711,7 @@ func (dp *DataPartition) fetchReplicasFromMaster() (isLeader bool, replicas []st
|
||||
time.Sleep(10 * time.Second)
|
||||
}
|
||||
|
||||
for _, host := range partition.Hosts {
|
||||
replicas = append(replicas, host)
|
||||
}
|
||||
replicas = append(replicas, partition.Hosts...)
|
||||
if partition.Hosts != nil && len(partition.Hosts) >= 1 {
|
||||
leaderAddr := strings.Split(partition.Hosts[0], ":")
|
||||
if len(leaderAddr) == 2 && strings.TrimSpace(leaderAddr[0]) == LocalIP {
|
||||
@ -1029,7 +951,6 @@ func (dp *DataPartition) getRepairConn(target string) (net.Conn, error) {
|
||||
func (dp *DataPartition) putRepairConn(conn net.Conn, forceClose bool) {
|
||||
log.LogDebugf("action[putRepairConn], forceClose: %v", forceClose)
|
||||
dp.dataNode.putRepairConnFunc(conn, forceClose)
|
||||
return
|
||||
}
|
||||
|
||||
func (dp *DataPartition) isNormalType() bool {
|
||||
|
||||
@ -160,14 +160,6 @@ func UnmarshalOldVersionRandWriteOpItem(raw []byte) (result *rndWrtOpItem, err e
|
||||
return
|
||||
}
|
||||
|
||||
func (dp *DataPartition) checkWriteErrs(errMsg string) (ignore bool) {
|
||||
// file has been deleted when applying the raft log
|
||||
if strings.Contains(errMsg, storage.ExtentHasBeenDeletedError.Error()) || strings.Contains(errMsg, storage.ExtentNotFoundError.Error()) {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// CheckLeader checks if itself is the leader during read
|
||||
func (dp *DataPartition) CheckLeader(request *repl.Packet, connect net.Conn) (err error) {
|
||||
// and use another getRaftLeaderAddr() to return the actual address
|
||||
@ -200,7 +192,6 @@ func (si *ItemIterator) ApplyIndex() uint64 {
|
||||
|
||||
// Close Closes the iterator.
|
||||
func (si *ItemIterator) Close() {
|
||||
return
|
||||
}
|
||||
|
||||
// Next returns the next item in the iterator.
|
||||
|
||||
@ -146,7 +146,6 @@ func (dp *DataPartition) stopRaft() {
|
||||
log.LogErrorf("[FATAL] stop raft partition(%v)", dp.partitionID)
|
||||
dp.raftPartition.Stop()
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (dp *DataPartition) CanRemoveRaftMember(peer proto.Peer, force bool) error {
|
||||
@ -725,12 +724,6 @@ func (dp *DataPartition) getLeaderMaxExtentIDAndPartitionSize() (maxExtentID, Pa
|
||||
return dp.getMaxExtentIDAndPartitionSize(target)
|
||||
}
|
||||
|
||||
// Get the MaxExtentID partition from the leader.
|
||||
func (dp *DataPartition) getMemberExtentIDAndPartitionSize() (maxExtentID, PartitionSize uint64, err error) {
|
||||
target := dp.getReplicaAddr(1)
|
||||
return dp.getMaxExtentIDAndPartitionSize(target)
|
||||
}
|
||||
|
||||
func (dp *DataPartition) broadcastMinAppliedID(minAppliedID uint64) (err error) {
|
||||
for i := 0; i < dp.getReplicaLen(); i++ {
|
||||
p := NewPacketToBroadcastMinAppliedID(dp.partitionID, minAppliedID)
|
||||
@ -867,6 +860,4 @@ func (dp *DataPartition) updateMaxMinAppliedID() {
|
||||
log.LogDebugf("[updateMaxMinAppliedID] PartitionID(%v) localID(%v) OK! oldMaxID(%v) newMaxID(%v)",
|
||||
dp.partitionID, dp.appliedID, dp.maxAppliedID, maxAppliedID)
|
||||
dp.maxAppliedID = maxAppliedID
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
@ -201,7 +201,7 @@ func doStart(server common.Server, cfg *config.Config) (err error) {
|
||||
return errors.New("Invalid node Type!")
|
||||
}
|
||||
|
||||
s.stopC = make(chan bool, 0)
|
||||
s.stopC = make(chan bool)
|
||||
|
||||
// parse the config file
|
||||
if err = s.parseConfig(cfg); err != nil {
|
||||
@ -288,7 +288,7 @@ func (s *DataNode) parseConfig(cfg *config.Config) (err error) {
|
||||
port = cfg.GetString(proto.ListenPort)
|
||||
s.bindIp = cfg.GetBool(proto.BindIpKey)
|
||||
serverPort = port
|
||||
if regexpPort, err = regexp.Compile("^(\\d)+$"); err != nil {
|
||||
if regexpPort, err = regexp.Compile(`^(\d)+$`); err != nil {
|
||||
return fmt.Errorf("Err:no port")
|
||||
}
|
||||
if !regexpPort.MatchString(port) {
|
||||
@ -404,12 +404,12 @@ func (s *DataNode) startSpaceManager(cfg *config.Config) (err error) {
|
||||
if s.clusterUuidEnable {
|
||||
if err = config.CheckOrStoreClusterUuid(path, s.clusterUuid, false); err != nil {
|
||||
log.LogErrorf("CheckOrStoreClusterUuid failed: %v", err)
|
||||
return errors.New(fmt.Sprintf("CheckOrStoreClusterUuid failed: %v", err.Error()))
|
||||
return fmt.Errorf("CheckOrStoreClusterUuid failed: %v", err.Error())
|
||||
}
|
||||
}
|
||||
reservedSpace, err := strconv.ParseUint(arr[1], 10, 64)
|
||||
if err != nil {
|
||||
return errors.New(fmt.Sprintf("Invalid disk reserved space. Error: %s", err.Error()))
|
||||
return fmt.Errorf("Invalid disk reserved space. Error: %s", err.Error())
|
||||
}
|
||||
|
||||
if reservedSpace < DefaultDiskRetainMin {
|
||||
@ -715,7 +715,6 @@ func (s *DataNode) serveSmuxConn(conn net.Conn) {
|
||||
}
|
||||
go s.serveSmuxStream(stream)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) serveSmuxStream(stream *smux.Stream) {
|
||||
@ -823,10 +822,8 @@ func (s *DataNode) initConnPool() {
|
||||
s.putRepairConnFunc = func(conn net.Conn, forceClose bool) {
|
||||
log.LogDebugf("[dataNode.putRepairConnFunc] put tcp conn, addr(%v), forceClose(%v)", conn.RemoteAddr().String(), forceClose)
|
||||
gConnPool.PutConnect(conn.(*net.TCPConn), forceClose)
|
||||
return
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) closeSmuxConnPool() {
|
||||
@ -834,7 +831,6 @@ func (s *DataNode) closeSmuxConnPool() {
|
||||
s.smuxConnPool.Close()
|
||||
log.LogDebugf("action[stopSmuxService] stop smux conn pool")
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) shallDegrade() bool {
|
||||
@ -846,10 +842,7 @@ func (s *DataNode) shallDegrade() bool {
|
||||
return false
|
||||
}
|
||||
cnt := atomic.LoadUint64(&s.metricsCnt)
|
||||
if cnt%uint64(level) == 0 {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
return cnt%uint64(level) != 0
|
||||
}
|
||||
|
||||
func (s *DataNode) scheduleTask() {
|
||||
@ -876,9 +869,7 @@ func (s *DataNode) startCpuSample() {
|
||||
func (s *DataNode) scheduleToCheckLackPartitions() {
|
||||
go func() {
|
||||
for {
|
||||
var err error
|
||||
lackPartitionsInMem := make([]uint64, 0)
|
||||
lackPartitionsInMem, err = s.checkLocalPartitionMatchWithMaster()
|
||||
lackPartitionsInMem, err := s.checkLocalPartitionMatchWithMaster()
|
||||
if err != nil {
|
||||
log.LogError(err)
|
||||
}
|
||||
@ -889,8 +880,7 @@ func (s *DataNode) scheduleToCheckLackPartitions() {
|
||||
}
|
||||
s.space.stats.updateMetricLackPartitionsInMem(uint64(len(lackPartitionsInMem)))
|
||||
|
||||
lackPartitionsInDisk := make([]uint64, 0)
|
||||
lackPartitionsInDisk = s.checkPartitionInMemoryMatchWithInDisk()
|
||||
lackPartitionsInDisk := s.checkPartitionInMemoryMatchWithInDisk()
|
||||
if len(lackPartitionsInDisk) > 0 {
|
||||
err = fmt.Errorf("action[scheduleToLackDataPartitions] lackPartitions %v in datanode %v disk",
|
||||
lackPartitionsInDisk, s.localServerAddr)
|
||||
@ -904,10 +894,7 @@ func (s *DataNode) scheduleToCheckLackPartitions() {
|
||||
}
|
||||
|
||||
func IsDiskErr(errMsg string) bool {
|
||||
if strings.Contains(errMsg, syscall.EIO.Error()) || strings.Contains(errMsg, syscall.EROFS.Error()) ||
|
||||
strings.Contains(errMsg, syscall.EACCES.Error()) {
|
||||
return true
|
||||
}
|
||||
|
||||
return false
|
||||
return strings.Contains(errMsg, syscall.EIO.Error()) ||
|
||||
strings.Contains(errMsg, syscall.EROFS.Error()) ||
|
||||
strings.Contains(errMsg, syscall.EACCES.Error())
|
||||
}
|
||||
|
||||
@ -246,7 +246,6 @@ func (s *DataNode) getExtentAPI(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
s.buildSuccessResp(w, extentInfo)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) getBlockCrcAPI(w http.ResponseWriter, r *http.Request) {
|
||||
@ -279,7 +278,6 @@ func (s *DataNode) getBlockCrcAPI(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
s.buildSuccessResp(w, blocks)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) getTinyDeleted(w http.ResponseWriter, r *http.Request) {
|
||||
@ -307,7 +305,6 @@ func (s *DataNode) getTinyDeleted(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
s.buildSuccessResp(w, extentInfo)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) getNormalDeleted(w http.ResponseWriter, r *http.Request) {
|
||||
@ -335,7 +332,6 @@ func (s *DataNode) getNormalDeleted(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
s.buildSuccessResp(w, extentInfo)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) setQosEnable() func(http.ResponseWriter, *http.Request) {
|
||||
@ -369,7 +365,6 @@ func (s *DataNode) getSmuxPoolStat() func(http.ResponseWriter, *http.Request) {
|
||||
}
|
||||
stat := s.smuxConnPool.GetStat()
|
||||
s.buildSuccessResp(w, stat)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
@ -415,7 +410,6 @@ func (s *DataNode) genClusterVersionFile(w http.ResponseWriter, r *http.Request)
|
||||
}
|
||||
}
|
||||
s.buildSuccessResp(w, "Generate cluster version file success")
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) buildSuccessResp(w http.ResponseWriter, data interface{}) {
|
||||
|
||||
@ -57,7 +57,7 @@ func NewSpaceManager(dataNode *DataNode) *SpaceManager {
|
||||
space.diskList = make([]string, 0)
|
||||
space.partitions = make(map[uint64]*DataPartition)
|
||||
space.stats = NewStats(dataNode.zoneName)
|
||||
space.stopC = make(chan bool, 0)
|
||||
space.stopC = make(chan bool)
|
||||
space.dataNode = dataNode
|
||||
space.diskUtils = make(map[string]*atomicutil.Float64)
|
||||
go space.statUpdateScheduler()
|
||||
|
||||
@ -217,8 +217,6 @@ func (s *DataNode) handlePacketToCreateExtent(p *repl.Packet) {
|
||||
partition.disk.allocCheckLimit(proto.IopsWriteType, 1)
|
||||
|
||||
err = partition.ExtentStore().Create(p.ExtentID)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpCreateDataPartition packet.
|
||||
@ -260,8 +258,6 @@ func (s *DataNode) handlePacketToCreateDataPartition(p *repl.Packet) {
|
||||
return
|
||||
}
|
||||
p.PacketOkWithBody([]byte(dp.Disk().Path))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) commitCreateVersion(volumeID string, verSeq uint64) (err error) {
|
||||
@ -652,8 +648,6 @@ func (s *DataNode) handleMarkDeletePacket(p *repl.Packet, c net.Conn) {
|
||||
partition.disk.allocCheckLimit(proto.IopsWriteType, 1)
|
||||
partition.ExtentStore().MarkDelete(p.ExtentID, 0, 0)
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpMarkDelete packet.
|
||||
@ -687,8 +681,6 @@ func (s *DataNode) handleBatchMarkDeletePacket(p *repl.Packet, c net.Conn) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpWrite packet.
|
||||
@ -784,7 +776,6 @@ func (s *DataNode) handleWritePacket(p *repl.Packet) {
|
||||
}
|
||||
}
|
||||
s.incDiskErrCnt(p.PartitionID, err, WriteFlag)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handleRandomWritePacket(p *repl.Packet) {
|
||||
@ -900,8 +891,6 @@ func (s *DataNode) handleStreamReadPacket(p *repl.Packet, connect net.Conn, isRe
|
||||
return
|
||||
}
|
||||
s.extentRepairReadPacket(p, connect, isRepairRead)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handleExtentRepairReadPacket(p *repl.Packet, connect net.Conn, isRepairRead bool) {
|
||||
@ -1002,8 +991,6 @@ func (s *DataNode) extentRepairReadPacket(p *repl.Packet, connect net.Conn, isRe
|
||||
log.LogReadf(logContent)
|
||||
}
|
||||
p.PacketOkReply()
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToGetAllWatermarks(p *repl.Packet) {
|
||||
@ -1027,9 +1014,12 @@ func (s *DataNode) handlePacketToGetAllWatermarks(p *repl.Packet) {
|
||||
p.PackErrorBody(ActionGetAllExtentWatermarks, err.Error())
|
||||
} else {
|
||||
buf, err = json.Marshal(fInfoList)
|
||||
p.PacketOkWithByte(buf)
|
||||
if err != nil {
|
||||
p.PackErrorBody(ActionGetAllExtentWatermarks, err.Error())
|
||||
} else {
|
||||
p.PacketOkWithByte(buf)
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) writeEmptyPacketOnTinyExtentRepairRead(reply *repl.Packet, newOffset, currentOffset int64, connect net.Conn) (replySize int64, err error) {
|
||||
@ -1137,7 +1127,6 @@ func (s *DataNode) tinyExtentRepairRead(request *repl.Packet, connect net.Conn)
|
||||
}
|
||||
|
||||
request.PacketOkReply()
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToReadTinyDeleteRecordFile(p *repl.Packet, connect net.Conn) {
|
||||
@ -1181,8 +1170,6 @@ func (s *DataNode) handlePacketToReadTinyDeleteRecordFile(p *repl.Packet, connec
|
||||
offset += int64(currReadSize)
|
||||
}
|
||||
p.PacketOkReply()
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpNotifyReplicasToRepair packet.
|
||||
@ -1197,7 +1184,6 @@ func (s *DataNode) handlePacketToNotifyExtentRepair(p *repl.Packet) {
|
||||
}
|
||||
partition.DoExtentStoreRepair(mf)
|
||||
p.PacketOkReply()
|
||||
return
|
||||
}
|
||||
|
||||
// Handle OpBroadcastMinAppliedID
|
||||
@ -1209,7 +1195,6 @@ func (s *DataNode) handleBroadcastMinAppliedID(p *repl.Packet) {
|
||||
}
|
||||
log.LogDebugf("[handleBroadcastMinAppliedID] partition(%v) minAppliedID(%v)", partition.partitionID, minAppliedID)
|
||||
p.PacketOkReply()
|
||||
return
|
||||
}
|
||||
|
||||
// Handle handlePacketToGetAppliedID packet.
|
||||
@ -1220,7 +1205,6 @@ func (s *DataNode) handlePacketToGetAppliedID(p *repl.Packet) {
|
||||
binary.BigEndian.PutUint64(buf, appliedID)
|
||||
p.PacketOkWithBody(buf)
|
||||
p.AddMesgLog(fmt.Sprintf("_AppliedID(%v)", appliedID))
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToGetPartitionSize(p *repl.Packet) {
|
||||
@ -1230,8 +1214,6 @@ func (s *DataNode) handlePacketToGetPartitionSize(p *repl.Packet) {
|
||||
binary.BigEndian.PutUint64(buf, uint64(usedSize))
|
||||
p.AddMesgLog(fmt.Sprintf("partitionSize_(%v)", usedSize))
|
||||
p.PacketOkWithBody(buf)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToGetMaxExtentIDAndPartitionSize(p *repl.Packet) {
|
||||
@ -1242,8 +1224,6 @@ func (s *DataNode) handlePacketToGetMaxExtentIDAndPartitionSize(p *repl.Packet)
|
||||
binary.BigEndian.PutUint64(buf[0:8], uint64(maxExtentID))
|
||||
binary.BigEndian.PutUint64(buf[8:16], totalPartitionSize)
|
||||
p.PacketOkWithBody(buf)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToDecommissionDataPartition(p *repl.Packet) {
|
||||
@ -1303,7 +1283,6 @@ func (s *DataNode) handlePacketToDecommissionDataPartition(p *repl.Packet) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) {
|
||||
@ -1364,8 +1343,6 @@ func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) {
|
||||
}
|
||||
}
|
||||
log.LogInfof("action[handlePacketToAddDataPartitionRaftMember] after ChangeRaftMember %v, partition id %v", req.AddPeer, &req.PartitionId)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) {
|
||||
@ -1452,7 +1429,6 @@ func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) {
|
||||
}
|
||||
log.LogDebugf("action[handlePacketToRemoveDataPartitionRaftMember] CanRemoveRaftMember complete "+
|
||||
"req %v dp %v ", p.GetReqID(), dp.partitionID)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) handlePacketToDataPartitionTryToLeader(p *repl.Packet) {
|
||||
@ -1484,7 +1460,6 @@ func (s *DataNode) handlePacketToDataPartitionTryToLeader(p *repl.Packet) {
|
||||
return
|
||||
}
|
||||
err = dp.raftPartition.TryToLeader(dp.partitionID)
|
||||
return
|
||||
}
|
||||
|
||||
func (s *DataNode) forwardToRaftLeader(dp *DataPartition, p *repl.Packet, force bool) (ok bool, err error) {
|
||||
|
||||
Loading…
Reference in New Issue
Block a user