mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
200 lines
6.1 KiB
Go
200 lines
6.1 KiB
Go
// Copyright 2023 The CubeFS Authors.
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
|
|
// implied. See the License for the specific language governing
|
|
// permissions and limitations under the License.
|
|
|
|
package lcnode
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/cubefs/cubefs/proto"
|
|
"github.com/cubefs/cubefs/util/auditlog"
|
|
"github.com/cubefs/cubefs/util/log"
|
|
)
|
|
|
|
func (l *LcNode) opMasterHeartbeat(conn net.Conn, p *proto.Packet, remoteAddr string) (err error) {
|
|
data := p.Data
|
|
responseAckOKToMaster(conn, p)
|
|
|
|
var (
|
|
req = &proto.HeartBeatRequest{}
|
|
resp = &proto.LcNodeHeartbeatResponse{
|
|
LcScanningTasks: make(map[string]*proto.LcNodeRuleTaskResponse),
|
|
SnapshotScanningTasks: make(map[string]*proto.SnapshotVerDelTaskResponse),
|
|
}
|
|
adminTask = &proto.AdminTask{
|
|
Request: req,
|
|
}
|
|
)
|
|
|
|
go func() {
|
|
start := time.Now()
|
|
decode := json.NewDecoder(bytes.NewBuffer(data))
|
|
decode.UseNumber()
|
|
if err = decode.Decode(adminTask); err != nil {
|
|
resp.Status = proto.TaskFailed
|
|
resp.Result = fmt.Sprintf("lcnode(%v) heartbeat decode err(%v)", l.localServerAddr, err.Error())
|
|
goto end
|
|
}
|
|
|
|
l.scannerMutex.RLock()
|
|
for _, scanner := range l.lcScanners {
|
|
result := &proto.LcNodeRuleTaskResponse{
|
|
ID: scanner.ID,
|
|
LcNode: l.localServerAddr,
|
|
StartTime: &scanner.now,
|
|
Volume: scanner.Volume,
|
|
RcvStop: scanner.receiveStop,
|
|
Rule: scanner.rule,
|
|
LcNodeRuleTaskStatistics: proto.LcNodeRuleTaskStatistics{
|
|
TotalFileScannedNum: atomic.LoadInt64(&scanner.currentStat.TotalFileScannedNum),
|
|
TotalFileExpiredNum: atomic.LoadInt64(&scanner.currentStat.TotalFileExpiredNum),
|
|
TotalDirScannedNum: atomic.LoadInt64(&scanner.currentStat.TotalDirScannedNum),
|
|
ExpiredDeleteNum: atomic.LoadInt64(&scanner.currentStat.ExpiredDeleteNum),
|
|
ExpiredMToHddNum: atomic.LoadInt64(&scanner.currentStat.ExpiredMToHddNum),
|
|
ExpiredMToBlobstoreNum: atomic.LoadInt64(&scanner.currentStat.ExpiredMToBlobstoreNum),
|
|
ExpiredMToHddBytes: atomic.LoadInt64(&scanner.currentStat.ExpiredMToHddBytes),
|
|
ExpiredMToBlobstoreBytes: atomic.LoadInt64(&scanner.currentStat.ExpiredMToBlobstoreBytes),
|
|
ExpiredSkipNum: atomic.LoadInt64(&scanner.currentStat.ExpiredSkipNum),
|
|
ErrorDeleteNum: atomic.LoadInt64(&scanner.currentStat.ErrorDeleteNum),
|
|
ErrorMToHddNum: atomic.LoadInt64(&scanner.currentStat.ErrorMToHddNum),
|
|
ErrorMToBlobstoreNum: atomic.LoadInt64(&scanner.currentStat.ErrorMToBlobstoreNum),
|
|
ErrorReadDirNum: atomic.LoadInt64(&scanner.currentStat.ErrorReadDirNum),
|
|
},
|
|
}
|
|
resp.LcScanningTasks[scanner.ID] = result
|
|
}
|
|
for _, scanner := range l.snapshotScanners {
|
|
info := &proto.SnapshotVerDelTaskResponse{
|
|
ID: scanner.ID,
|
|
LcNode: l.localServerAddr,
|
|
SnapshotVerDelTask: scanner.verDelReq.Task,
|
|
SnapshotStatistics: proto.SnapshotStatistics{
|
|
VolName: scanner.Volume,
|
|
VerSeq: scanner.getTaskVerSeq(),
|
|
TotalInodeNum: atomic.LoadInt64(&scanner.currentStat.TotalInodeNum),
|
|
FileNum: atomic.LoadInt64(&scanner.currentStat.FileNum),
|
|
DirNum: atomic.LoadInt64(&scanner.currentStat.DirNum),
|
|
ErrorSkippedNum: atomic.LoadInt64(&scanner.currentStat.ErrorSkippedNum),
|
|
},
|
|
}
|
|
resp.SnapshotScanningTasks[scanner.ID] = info
|
|
}
|
|
l.scannerMutex.RUnlock()
|
|
|
|
resp.LcTaskCountLimit = lcNodeTaskCountLimit
|
|
resp.Status = proto.TaskSucceeds
|
|
|
|
end:
|
|
adminTask.Response = resp
|
|
l.respondToMaster(adminTask)
|
|
msg := fmt.Sprintf("from(%v), adminTask(%+v), resp(%+v), %v", remoteAddr, adminTask, resp, time.Since(start).String())
|
|
log.LogInfof("MasterHeartbeat %v ", msg)
|
|
auditlog.LogMasterOp("MasterHeartbeat", msg, err)
|
|
}()
|
|
|
|
l.lastHeartbeat = time.Now()
|
|
log.LogDebugf("lastHeartbeat: %v", l.lastHeartbeat)
|
|
return
|
|
}
|
|
|
|
func (l *LcNode) opLcScan(conn net.Conn, p *proto.Packet) (err error) {
|
|
data := p.Data
|
|
|
|
responseAckOKToMaster(conn, p)
|
|
|
|
go func() {
|
|
var (
|
|
req = &proto.LcNodeRuleTaskRequest{}
|
|
resp = &proto.LcNodeRuleTaskResponse{}
|
|
adminTask = &proto.AdminTask{
|
|
Request: req,
|
|
}
|
|
)
|
|
|
|
decoder := json.NewDecoder(bytes.NewBuffer(data))
|
|
decoder.UseNumber()
|
|
if err = decoder.Decode(adminTask); err != nil {
|
|
resp.LcNode = l.localServerAddr
|
|
resp.Status = proto.TaskFailed
|
|
resp.Done = true
|
|
resp.StartErr = err.Error()
|
|
adminTask.Response = resp
|
|
l.respondToMaster(adminTask)
|
|
return
|
|
}
|
|
|
|
l.startLcScan(adminTask)
|
|
l.respondToMaster(adminTask)
|
|
}()
|
|
|
|
return
|
|
}
|
|
|
|
func (l *LcNode) respondToMaster(task *proto.AdminTask) {
|
|
// handle panic
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
log.LogErrorf("respondToMaster err: %v", r)
|
|
}
|
|
}()
|
|
if err := l.mc.NodeAPI().ResponseLcNodeTask(task); err != nil {
|
|
log.LogErrorf("respondToMaster err: %v, task: %v", err, task)
|
|
}
|
|
}
|
|
|
|
func (l *LcNode) opSnapshotVerDel(conn net.Conn, p *proto.Packet) (err error) {
|
|
data := p.Data
|
|
|
|
responseAckOKToMaster(conn, p)
|
|
|
|
go func() {
|
|
var (
|
|
req = &proto.SnapshotVerDelTaskRequest{}
|
|
resp = &proto.SnapshotVerDelTaskResponse{}
|
|
adminTask = &proto.AdminTask{
|
|
Request: req,
|
|
}
|
|
)
|
|
|
|
decoder := json.NewDecoder(bytes.NewBuffer(data))
|
|
decoder.UseNumber()
|
|
if err = decoder.Decode(adminTask); err != nil {
|
|
resp.Status = proto.TaskFailed
|
|
resp.Result = err.Error()
|
|
adminTask.Response = resp
|
|
l.respondToMaster(adminTask)
|
|
return
|
|
}
|
|
|
|
l.startSnapshotScan(adminTask)
|
|
l.respondToMaster(adminTask)
|
|
}()
|
|
|
|
return
|
|
}
|
|
|
|
func responseAckOKToMaster(conn net.Conn, p *proto.Packet) {
|
|
go func() {
|
|
p.PacketOkReply()
|
|
if err := p.WriteToConn(conn); err != nil {
|
|
log.LogErrorf("ack master response: %s", err.Error())
|
|
}
|
|
}()
|
|
}
|