From 984f3daf1160b0bbf307bb04d4e3a825c67c66f7 Mon Sep 17 00:00:00 2001 From: clinx Date: Tue, 12 Aug 2025 10:59:31 +0800 Subject: [PATCH] fix(client): Add retry logic when an error occurs while fetching remote configuration information with: #1000275887 Signed-off-by: clinx (cherry picked from commit dff1e3525b43790a656d411cc6d52f47e3f38a93) --- .../flashgroupmanager/admin_task_manager.go | 31 ++++++ remotecache/flashgroupmanager/cluster.go | 17 +++ remotecache/flashgroupmanager/flash_node.go | 10 ++ remotecache/flashgroupmanager/http_server.go | 104 +++++++++++++++++- sdk/remotecache/client.go | 24 +++- 5 files changed, 179 insertions(+), 7 deletions(-) diff --git a/remotecache/flashgroupmanager/admin_task_manager.go b/remotecache/flashgroupmanager/admin_task_manager.go index b6057a982..a8be9fe86 100644 --- a/remotecache/flashgroupmanager/admin_task_manager.go +++ b/remotecache/flashgroupmanager/admin_task_manager.go @@ -244,3 +244,34 @@ func (sender *AdminTaskManager) AddTask(t *proto.AdminTask) { sender.TaskMap[t.ID] = t } } + +func (sender *AdminTaskManager) syncSendAdminTask(task *proto.AdminTask) (packet *proto.Packet, err error) { + packet, err = sender.buildPacket(task) + if err != nil { + return nil, errors.Trace(err, "action[syncSendAdminTask build packet failed,task:%v]", task.ID) + } + log.LogInfof("action[syncSendAdminTask],task[%s], op %s, reqId %d", task.ToString(), packet.GetOpMsg(), packet.GetReqID()) + conn, err := sender.getConn() + if err != nil { + return nil, errors.Trace(err, "action[syncSendAdminTask get conn failed,task:%v]", task.ID) + } + defer func() { + if err == nil { + sender.putConn(conn, false) + } else { + sender.putConn(conn, true) + } + }() + if err = packet.WriteToConn(conn); err != nil { + return nil, errors.Trace(err, "action[syncSendAdminTask],WriteToConn failed,task:%v,reqID[%v]", task.ID, packet.ReqID) + } + if err = packet.ReadFromConnWithVer(conn, proto.SyncSendTaskDeadlineTime); err != nil { + return nil, errors.Trace(err, "action[syncSendAdminTask],ReadFromConn failed task:%v,reqID[%v]", task.ID, packet.ReqID) + } + if packet.ResultCode != proto.OpOk { + err = fmt.Errorf("result code[%v],msg[%v]", packet.ResultCode, string(packet.Data)) + log.LogErrorf("action[syncSendAdminTask],task:%v,reqID[%v],err[%v],", task.ID, packet.ReqID, err) + return + } + return packet, nil +} diff --git a/remotecache/flashgroupmanager/cluster.go b/remotecache/flashgroupmanager/cluster.go index 0e402f30a..d50e70f8a 100644 --- a/remotecache/flashgroupmanager/cluster.go +++ b/remotecache/flashgroupmanager/cluster.go @@ -725,3 +725,20 @@ func (c *Cluster) syncUpdateFlashNode(flashNode *FlashNode) (err error) { func (c *Cluster) tryToChangeLeaderByHost() error { return c.partition.TryToLeader(1) } + +func (c *Cluster) syncFlashNodeSetIOLimitTasks(tasks []*proto.AdminTask) { + for _, t := range tasks { + if t == nil { + continue + } + node, err := c.peekFlashNode(t.OperatorAddr) + if err != nil { + log.LogWarn(fmt.Sprintf("action[syncFlashNodeHeartbeatTasks],nodeAddr:%v,taskID:%v,err:%v", t.OperatorAddr, t.ID, err.Error())) + continue + } + if _, err = node.TaskManager.syncSendAdminTask(t); err != nil { + log.LogWarn(fmt.Sprintf("action[syncFlashNodeHeartbeatTasks],nodeAddr:%v,taskID:%v,err:%v", t.OperatorAddr, t.ID, err.Error())) + continue + } + } +} diff --git a/remotecache/flashgroupmanager/flash_node.go b/remotecache/flashgroupmanager/flash_node.go index c7bee98a4..8deb94202 100644 --- a/remotecache/flashgroupmanager/flash_node.go +++ b/remotecache/flashgroupmanager/flash_node.go @@ -129,3 +129,13 @@ func (flashNode *FlashNode) createHeartbeatTask(masterAddr string, flashNodeHand task = proto.NewAdminTask(proto.OpFlashNodeHeartbeat, flashNode.Addr, request) return } + +func (flashNode *FlashNode) createSetIOLimitsTask(flow, iocc, factor int, opCode uint8) (task *proto.AdminTask) { + request := &proto.FlashNodeSetIOLimitsRequest{ + Flow: flow, + Iocc: iocc, + Factor: factor, + } + task = proto.NewAdminTask(opCode, flashNode.Addr, request) + return +} diff --git a/remotecache/flashgroupmanager/http_server.go b/remotecache/flashgroupmanager/http_server.go index ade401e4e..ed525aa86 100644 --- a/remotecache/flashgroupmanager/http_server.go +++ b/remotecache/flashgroupmanager/http_server.go @@ -6,6 +6,7 @@ import ( "net/http/httputil" "time" + "github.com/cubefs/cubefs/cmd/common" "github.com/cubefs/cubefs/proto" "github.com/cubefs/cubefs/util/config" "github.com/cubefs/cubefs/util/exporter" @@ -91,7 +92,8 @@ func (m *FlashGroupManager) registerAPIRoutes(router *mux.Router) { router.NewRoute().Methods(http.MethodGet, http.MethodPost).Path(proto.FlashNodeRemove).HandlerFunc(m.removeFlashNode) router.NewRoute().Methods(http.MethodGet, http.MethodPost).Path(proto.FlashNodeRemoveAllInactive).HandlerFunc(m.removeAllInactiveFlashNodes) router.NewRoute().Methods(http.MethodGet).Path(proto.FlashNodeGet).HandlerFunc(m.getFlashNode) - + router.NewRoute().Methods(http.MethodGet, http.MethodPost).Path(proto.FlashNodeSetReadIOLimits).HandlerFunc(m.setFlashNodeReadIOLimits) + router.NewRoute().Methods(http.MethodGet, http.MethodPost).Path(proto.FlashNodeSetWriteIOLimits).HandlerFunc(m.setFlashNodeWriteIOLimits) router.NewRoute().Methods(http.MethodGet, http.MethodPost). Path(proto.GetFlashNodeTaskResponse). HandlerFunc(m.handleFlashNodeTaskResponse) @@ -178,3 +180,103 @@ func (m *FlashGroupManager) newReverseProxy() *httputil.ReverseProxy { func (m *FlashGroupManager) proxy(w http.ResponseWriter, r *http.Request) { m.reverseProxy.ServeHTTP(w, r) } + +func (m *FlashGroupManager) setFlashNodeReadIOLimits(w http.ResponseWriter, r *http.Request) { + var ( + flow common.Int + iocc common.Int + factor common.Int + readFlow int64 + readIocc int64 + readFactor int64 + err error + ) + + if err = parseArgs(r, flow.Flow().OmitEmpty().OnEmpty(func() error { + readFlow = -1 + return nil + }).OnValue(func() error { + readFlow = flow.V + return nil + }), + iocc.Iocc().OmitEmpty().OnEmpty(func() error { + readIocc = -1 + return nil + }).OnValue(func() error { + readIocc = iocc.V + return nil + }), + factor.Factor().OmitEmpty().OnEmpty(func() error { + readFactor = -1 + return nil + }).OnValue(func() error { + readFactor = factor.V + return nil + })); err != nil { + sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) + return + } + log.LogDebugf("action[setFlashNodeReadIOLimits],flow[%v] iocc[%v] factor [%v]", + readFlow, readIocc, readFactor) + tasks := make([]*proto.AdminTask, 0) + m.cluster.flashNodeTopo.flashNodeMap.Range(func(key, value interface{}) bool { + flashNode := value.(*FlashNode) + if flashNode.isActiveAndEnable() { + task := flashNode.createSetIOLimitsTask(int(readFlow), int(readIocc), int(readFactor), proto.OpFlashNodeSetReadIOLimits) + tasks = append(tasks, task) + } + return true + }) + go m.cluster.syncFlashNodeSetIOLimitTasks(tasks) + sendOkReply(w, r, newSuccessHTTPReply("set ReadIOLimits for FlashNode is submit,check it later.")) +} + +func (m *FlashGroupManager) setFlashNodeWriteIOLimits(w http.ResponseWriter, r *http.Request) { + var ( + flow common.Int + iocc common.Int + factor common.Int + writeFlow int64 + writeIocc int64 + writeFactor int64 + err error + ) + + if err = parseArgs(r, flow.Flow().OmitEmpty().OnEmpty(func() error { + writeFlow = -1 + return nil + }).OnValue(func() error { + writeFlow = flow.V + return nil + }), + iocc.Iocc().OmitEmpty().OnEmpty(func() error { + writeIocc = -1 + return nil + }).OnValue(func() error { + writeIocc = iocc.V + return nil + }), + factor.Factor().OmitEmpty().OnEmpty(func() error { + writeFactor = -1 + return nil + }).OnValue(func() error { + writeFactor = factor.V + return nil + })); err != nil { + sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) + return + } + log.LogDebugf("action[setFlashNodeWriteIOLimits],flow[%v] iocc[%v] factor [%v]", + writeFlow, writeIocc, writeFactor) + tasks := make([]*proto.AdminTask, 0) + m.cluster.flashNodeTopo.flashNodeMap.Range(func(key, value interface{}) bool { + flashNode := value.(*FlashNode) + if flashNode.isActiveAndEnable() { + task := flashNode.createSetIOLimitsTask(int(writeFlow), int(writeIocc), int(writeFactor), proto.OpFlashNodeSetWriteIOLimits) + tasks = append(tasks, task) + } + return true + }) + go m.cluster.syncFlashNodeSetIOLimitTasks(tasks) + sendOkReply(w, r, newSuccessHTTPReply("set WriteIOLimits for FlashNode is submit,check it later.")) +} diff --git a/sdk/remotecache/client.go b/sdk/remotecache/client.go index c11c39a1c..0b51e7832 100755 --- a/sdk/remotecache/client.go +++ b/sdk/remotecache/client.go @@ -243,9 +243,15 @@ func (rc *RemoteCacheClient) UpdateFlashGroups() (err error) { fgv proto.FlashGroupView newFlashGroups = btree.New(32) ) - if fgv, err = rc.mc.AdminAPI().ClientFlashGroups(); err != nil { - log.LogWarnf("updateFlashGroups: err(%v)", err) - return + + for i := 0; i < 3; i++ { + if fgv, err = rc.mc.AdminAPI().ClientFlashGroups(); err == nil { + break + } + log.LogWarnf("updateFlashGroups: attempt %d failed, err(%v)", i+1, err) + if i == 2 { // Last attempt + return + } } log.LogDebugf("updateFlashGroups. get flashGroupView [%v]", fgv) rc.SetClusterEnable(fgv.Enable && len(fgv.FlashGroups) != 0) @@ -951,9 +957,15 @@ func (rc *RemoteCacheClient) Get(ctx context.Context, reqId, key string, from, t func (rc *RemoteCacheClient) updateRemoteCacheConfig() (err error) { var config *proto.RemoteCacheConfig - if config, err = rc.mc.AdminAPI().GetRemoteCacheConfig(); err != nil { - log.LogWarnf("updateRemoteCacheConfig: GetRemoteCacheConfig fail err(%v)", err) - return + + for i := 0; i < 3; i++ { + if config, err = rc.mc.AdminAPI().GetRemoteCacheConfig(); err == nil { + break + } + log.LogWarnf("updateRemoteCacheConfig: attempt %d failed, GetRemoteCacheConfig fail err(%v)", i+1, err) + if i == 2 { // Last attempt + return + } } log.LogInfof("updateRemoteCacheConfig: config(%v)", config)