From cd086fcb46eb629644f005f1fa37e1725dd32dd7 Mon Sep 17 00:00:00 2001 From: clinx Date: Mon, 20 Oct 2025 20:13:39 +0800 Subject: [PATCH] fix(flash): warm up task set force remote cache true with: #1000390069 Signed-off-by: clinx (cherry picked from commit 272660d7cb2e6a5a0b9ca6c62825a66d9754b4e4) --- master/flash_manual_task.go | 4 +-- master/flash_node.go | 4 ++- proto/distributed_cache.go | 3 ++ remotecache/flashnode/flashnode_op.go | 1 + sdk/data/stream/extent_client.go | 3 ++ sdk/remotecache/client.go | 52 +++++++++++++++++++-------- sdk/remotecache/read_task_queue.go | 6 ++-- 7 files changed, 53 insertions(+), 20 deletions(-) diff --git a/master/flash_manual_task.go b/master/flash_manual_task.go index d83a9f495..195320d8c 100644 --- a/master/flash_manual_task.go +++ b/master/flash_manual_task.go @@ -12,7 +12,7 @@ import ( ) const ( - minRequiredTTLSeconds = 30 * 60 // 30 minute + minRequiredTTLSeconds = 10 * 60 // 10 minute FlashTaskInteractiveTimeout = time.Minute * time.Duration(5) ScanManualInterval = time.Minute * time.Duration(2) CheckInteractiveInterval = time.Minute * time.Duration(2) @@ -50,7 +50,7 @@ func checkManualConfig(flt *proto.FlashManualTask, vol *Vol, fltMgr *flashManual switch flt.Action { case proto.FlashManualWarmupAction: if vol.remoteCacheTTL < minRequiredTTLSeconds && vol.remoteCacheTTL != 0 { - return fmt.Errorf("cache ttl %v of vol[%v] is too short and warm up can not work", vol.remoteCacheTTL, vol.Name) + return fmt.Errorf("cache ttl [%v]s of vol[%v] less than [%v]s and warm up can not work", vol.remoteCacheTTL, vol.Name, minRequiredTTLSeconds) } if err = isDuplicated(flt, fltMgr); err != nil { return err diff --git a/master/flash_node.go b/master/flash_node.go index 1d1a9949a..74b89fbac 100644 --- a/master/flash_node.go +++ b/master/flash_node.go @@ -278,7 +278,9 @@ func (m *Server) createFlashNodeManualTask(w http.ResponseWriter, r *http.Reques sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) return } - + if req.ManualTaskStatistics == nil { + req.ManualTaskStatistics = &proto.ManualTaskStatistics{} + } // Validate file size limits if req.ManualTaskConfig.MinFileSizeLimit > req.ManualTaskConfig.MaxFileSizeLimit { err = fmt.Errorf("MinFileSizeLimit(%d) cannot be greater than MaxFileSizeLimit(%d)", diff --git a/proto/distributed_cache.go b/proto/distributed_cache.go index 6175d5380..c9e7a7a36 100644 --- a/proto/distributed_cache.go +++ b/proto/distributed_cache.go @@ -580,6 +580,9 @@ func (flt *FlashManualTask) GetPathPrefix() string { func (flt *FlashManualTask) SetResponse(taskRsp *FlashNodeManualTaskResponse) { t := time.Now() flt.UpdateTime = &t + if flt.ManualTaskStatistics == nil { + flt.ManualTaskStatistics = &ManualTaskStatistics{} + } flt.ManualTaskStatistics.FlashNode = taskRsp.FlashNode flt.ManualTaskStatistics.TotalFileScannedNum = taskRsp.TotalFileScannedNum flt.ManualTaskStatistics.TotalFileCachedNum = taskRsp.TotalFileCachedNum diff --git a/remotecache/flashnode/flashnode_op.go b/remotecache/flashnode/flashnode_op.go index e4c8d6c5e..5f1ead786 100644 --- a/remotecache/flashnode/flashnode_op.go +++ b/remotecache/flashnode/flashnode_op.go @@ -844,6 +844,7 @@ func (f *FlashNode) getValidViewInfo(req *proto.FlashNodeManualTaskRequest) (met MetaWrapper: metaWrapper, NeedRemoteCache: true, HeartBeatPing: true, + ForceRemoteCache: true, } log.LogInfof("[NewS3Scanner] extentConfig: vol(%v) volStorageClass(%v) allowedStorageClass(%v), followerRead(%v)", extentConfig.Volume, extentConfig.VolStorageClass, extentConfig.VolAllowedStorageClass, extentConfig.FollowerRead) diff --git a/sdk/data/stream/extent_client.go b/sdk/data/stream/extent_client.go index 2e884fc0c..ddafeab10 100644 --- a/sdk/data/stream/extent_client.go +++ b/sdk/data/stream/extent_client.go @@ -927,6 +927,9 @@ func (client *ExtentClient) Close() error { _ = client.EvictStream(inode) } client.dataWrapper.Stop() + if client.RemoteCache.Started { + client.RemoteCache.Stop() + } return nil } diff --git a/sdk/remotecache/client.go b/sdk/remotecache/client.go index 0769aa65e..7e54e7245 100755 --- a/sdk/remotecache/client.go +++ b/sdk/remotecache/client.go @@ -117,6 +117,15 @@ type RemoteCacheClient struct { WriteChunkSize int64 // bytes } +func (rc *RemoteCacheClient) EnqueueConnTask(task *ConnPutTask) { + select { + case <-rc.stopC: + rc.processConnPut(task) + default: + rc.connPutChan <- task + } +} + func (as *AddressPingStats) Add(duration time.Duration) { as.Lock() defer as.Unlock() @@ -268,7 +277,22 @@ func (rc *RemoteCacheClient) consumeConnPut() { for { select { case <-rc.stopC: - log.LogDebugf("consumeConnPut: stopping due to stopC signal") + log.LogDebugf("consumeConnPut: stopC received, draining connPutChan for up to 10s") + func() { + timer := time.NewTimer(10 * time.Second) + defer timer.Stop() + for { + select { + case connTask := <-rc.connPutChan: + if connTask != nil { + rc.processConnPut(connTask) + } + case <-timer.C: + log.LogDebugf("consumeConnPut: drain window elapsed, exiting") + return + } + } + }() return case connTask := <-rc.connPutChan: if connTask != nil { @@ -584,7 +608,7 @@ func (rc *RemoteCacheClient) HeartBeat(addr string) (duration time.Duration, err packet.Opcode = proto.OpFlashSDKHeartbeat defer func() { - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) }() if conn, err = rc.conns.GetConnect(addr); err != nil { @@ -750,7 +774,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re bgTime := stat.BeginStat() forceClose := false defer func() { - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: forceClose} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: forceClose}) parts := strings.Split(addr, ":") if len(parts) > 0 && addr != "" { stat.EndStat(fmt.Sprintf("flashPutBlock:%v", parts[0]), err, bgTime, 1) @@ -957,7 +981,7 @@ func (rc *RemoteCacheClient) deleteRemoteBlock(key string, addr string) { return } defer func() { - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) }() p := proto.NewPacketReqID() p.Data = ([]byte)(key) @@ -1034,7 +1058,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int bgTime := stat.BeginStat() defer func() { forceClose := err != nil && !proto.IsFlashNodeLimitError(err) - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: forceClose} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: forceClose}) if err != nil && strings.Contains(err.Error(), "timeout") { err = proto.ErrorReadTimeout } @@ -1068,7 +1092,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int if err = reqPacket.WriteToConn(conn); err != nil { log.LogWarnf("FlashGroup Read: failed to write to addr(%v) err(%v) remoteCacheMultiRead(%v)", addr, err, rc.RemoteCacheMultiRead) - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) moved = fg.moveToUnknownRank(addr, err, rc.FlashNodeTimeoutCount) if rc.RemoteCacheMultiRead { log.LogInfof("Retrying due to write to addr(%v) failure err(%v)", addr, err) @@ -1082,7 +1106,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int break } log.LogWarnf("FlashGroup Read: getReadReply from addr(%v) err(%v) remoteCacheMultiRead(%v)", addr, err, rc.RemoteCacheMultiRead) - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) moved = fg.moveToUnknownRank(addr, err, rc.FlashNodeTimeoutCount) if rc.RemoteCacheMultiRead { log.LogInfof("Retrying due to getReadReply from addr(%v) failure err(%v)", addr, err) @@ -1123,7 +1147,7 @@ func (rc *RemoteCacheClient) Prepare(ctx context.Context, fg *FlashGroup, req *p return } defer func() { - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) }() if err = reqPacket.WriteToConn(conn); err != nil { @@ -1459,7 +1483,7 @@ func (reader *RemoteCacheReader) read(p []byte) (n int, err error) { atomic.AddInt64(&reader.alreadyReadLen, expectLen) if atomic.LoadInt64(&reader.alreadyReadLen) >= reader.needReadLen { reader.loadAll = true - reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: false} + reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: false}) reader.conn = nil } return int(expectLen), nil @@ -1469,7 +1493,7 @@ func (reader *RemoteCacheReader) Read(p []byte) (n int, err error) { defer func() { if err != nil && err != io.EOF && !reader.closed { reader.closed = true - reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: true} + reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: true}) if !reader.directRead { bytespool.Free(reader.buffer) } @@ -1535,7 +1559,7 @@ func (reader *RemoteCacheReader) Close() error { if !reader.closed { log.LogDebugf("RemoteCacheReader Close, reqId(%v)", reader.reqID) reader.closed = true - reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: reader.currOffset < reader.endOffset} + reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: reader.currOffset < reader.endOffset}) if !reader.directRead { bytespool.Free(reader.buffer) } @@ -1569,7 +1593,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool if err = op.ReqPacket.WriteToNoDeadLineConn(conn); err != nil { log.LogWarnf("%v FlashGroup Read: failed to write to addr(%v) err(%v)", op.LogPrefix, op.Addr, err) - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil}) if isLast && atomic.CompareAndSwapInt32(&op.HasResult, 0, 1) { op.Conn = nil op.BlockDataSize = 0 @@ -1591,7 +1615,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool log.LogDebugf("%v ReadObject getReadObjectReply from(%v) failed, reply(%v) error(%v)", op.LogPrefix, conn.RemoteAddr(), reply, err) } - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: !openConn} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: !openConn}) if openConn && atomic.CompareAndSwapInt32(&op.HasResult, 0, 1) { op.Conn = nil op.BlockDataSize = 0 @@ -1606,7 +1630,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool op.Err = nil atomic.StoreInt32(&op.EndLoop, 1) } else { - rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true} + rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true}) } } diff --git a/sdk/remotecache/read_task_queue.go b/sdk/remotecache/read_task_queue.go index 5564a140b..739651eed 100644 --- a/sdk/remotecache/read_task_queue.go +++ b/sdk/remotecache/read_task_queue.go @@ -220,20 +220,20 @@ func (rq *ReadTaskQueue) batchSmallObjectRead(taskWithPacket *ReadTaskWithPacket } conn.SetWriteDeadline(time.Now().Add(time.Duration(connTimeOut) * time.Millisecond)) if err = taskWithPacket.packet.WriteToNoDeadLineConn(conn); err != nil { - rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true} + rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true}) log.LogWarnf("batchSmallObjectRead: failed to write, firstPacktTime(%v) addr(%v) err(%v)", rq.rc.firstPacketTimeout, rq.flashAddr, err) return } replyPacket = NewFlashCacheReply() if err = replyPacket.ReadFromConnExt(conn, connTimeOut); err != nil { - rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true} + rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true}) log.LogWarnf("batchSmallObjectRead: failed to read, firstPacktTime(%v) addr(%v) err(%v)", rq.rc.firstPacketTimeout, rq.flashAddr, err) return } getResult = true atomic.StoreInt32(&taskWithPacket.finish, 1) - rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: false} + rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: false}) if replyPacket.ResultCode != proto.OpOk { err = fmt.Errorf(string(replyPacket.Data)) log.LogWarnf("batchSmallObjectRead: ResultCode NOK, replyPacket(%v), addr(%v), ResultCode(%v) err(%v)", replyPacket, rq.flashAddr, replyPacket.ResultCode, err.Error())