From d760b25ce111b7520387bd59b4e16da745ecb7a5 Mon Sep 17 00:00:00 2001 From: clinx Date: Mon, 23 Jun 2025 14:07:08 +0800 Subject: [PATCH] fix(flashnode): rate limiting based on actual written data flow close:#1000150848 Signed-off-by: clinx (cherry picked from commit 4320c3fc7ca3ed4665d656407ac008010fc511c5) --- .../flashnode/cachengine/block_test.go | 5 + remotecache/flashnode/flashnode_op.go | 93 ++++++++++--------- remotecache/flashnode/flashnode_op_test.go | 9 +- sdk/remotecache/client.go | 86 ++++++++++++----- util/unit.go | 5 +- 5 files changed, 122 insertions(+), 76 deletions(-) diff --git a/remotecache/flashnode/cachengine/block_test.go b/remotecache/flashnode/cachengine/block_test.go index 783017302..4830ef9e5 100644 --- a/remotecache/flashnode/cachengine/block_test.go +++ b/remotecache/flashnode/cachengine/block_test.go @@ -414,8 +414,13 @@ func testWriteSingleFileV2(t *testing.T) { file, _ := cacheBlock.GetOrOpenFileHandler() _, _ = file.ReadAt(crcBuf[:4], 8192+HeaderSize) require.Equal(t, crcSum1, binary.BigEndian.Uint32(crcBuf[:4])) + fdata := make([]byte, 4096) + _, _ = file.ReadAt(fdata[:4096], HeaderSize) + require.Equal(t, crcSum1, crc32.ChecksumIEEE(fdata[:4096])) _, _ = file.ReadAt(crcBuf[:4], 8192+HeaderSize+4) require.Equal(t, crcSum2, binary.BigEndian.Uint32(crcBuf[:4])) + _, _ = file.ReadAt(fdata[:4096], HeaderSize+4096) + require.Equal(t, crcSum2, crc32.ChecksumIEEE(fdata[:4096])) } func testWriteSingleFileErrorV2(t *testing.T) { diff --git a/remotecache/flashnode/flashnode_op.go b/remotecache/flashnode/flashnode_op.go index 1bc814273..b04752b3e 100644 --- a/remotecache/flashnode/flashnode_op.go +++ b/remotecache/flashnode/flashnode_op.go @@ -320,75 +320,82 @@ func (f *FlashNode) opCachePutBlock(conn net.Conn, p *proto.Packet) (err error) log.LogDebug("action[opCachePutBlock] create block key:"+uniKey+" logMsg:%s", p.LogMessage(p.GetOpMsg(), conn.RemoteAddr().String(), p.StartT, err)) } - missTaskDone := make(chan struct{}) allocSize := cachengine.CalcAllocSizeV2(int(req.BlockLen)) - err = f.limitWrite.TryRunAsync(context.Background(), allocSize, false, func() { - defer func() { - close(missTaskDone) - }() - if cb, err2, created := f.cacheEngine.CreateBlockV2(pDir, uniKey, req.TTL, uint32(allocSize), conn.RemoteAddr().String()); err2 != nil || created { - err = fmt.Errorf("already create block(%v)", cachengine.GenCacheBlockKeyV2(pDir, uniKey)) + if cb, err2, created := f.cacheEngine.CreateBlockV2(pDir, uniKey, req.TTL, uint32(allocSize), conn.RemoteAddr().String()); err2 != nil || created { + err = fmt.Errorf("already create block(%v)", cachengine.GenCacheBlockKeyV2(pDir, uniKey)) + return + } else { + p.PacketOkReply() + if err1 = p.WriteToConn(conn); err1 != nil { + log.LogErrorf("action[opCachePutBlock] write to conn %v", err1) return - } else { - p.PacketOkReply() - if e := p.WriteToConn(conn); e != nil { - log.LogErrorf("action[opCachePutBlock] write to conn %v", e) - } - bgTime1 := stat.BeginStat() - defer func() { - stat.EndStat("CachePutBlock:Write", err1, bgTime1, 1) - }() - var ( - totalWritten int64 - n int - readSize int - ) - buf := bytespool.Alloc(proto.PageSize) - defer bytespool.Free(buf) - crcBuf := bytespool.Alloc(cachengine.CRCLen) - defer bytespool.Free(crcBuf) - for { - readSize = proto.PageSize + } + bgTime1 := stat.BeginStat() + defer func() { + stat.EndStat("CachePutBlock:Write", err1, bgTime1, 1) + }() + var ( + totalWritten int64 + n int + readSize int + ) + writeLen := proto.PageSize + 4 + buf := bytespool.Alloc(writeLen) + defer bytespool.Free(buf) + for { + readSize = proto.PageSize + missTaskDone := make(chan struct{}) + err = f.limitWrite.TryRunAsync(context.Background(), writeLen, false, func() { + defer func() { + close(missTaskDone) + }() if totalWritten+int64(readSize) > req.BlockLen { readSize = int(req.BlockLen - totalWritten) } - if n, err1 = io.ReadFull(conn, buf[:proto.PageSize]); err1 != nil { + if n, err1 = io.ReadFull(conn, buf[:writeLen]); err1 != nil { + log.LogWarnf("action[opCachePutBlock] read data and crc from conn %v", err1) return } - if n != proto.PageSize { + if n != writeLen { err1 = syscall.EBADMSG - } - if n, err1 = io.ReadFull(conn, crcBuf[:cachengine.CRCLen]); err1 != nil { return } - if n != cachengine.CRCLen { - err1 = syscall.EBADMSG - } err1 = cb.WriteAtV2(&proto.FlashWriteParam{ Offset: totalWritten, Size: req.BlockLen, Data: buf[:proto.PageSize], - Crc: crcBuf[:cachengine.CRCLen], + Crc: buf[proto.PageSize:writeLen], DataSize: int64(readSize), }) + if err1 != nil { + return + } cachengine.UpdateWriteBytesMetric(proto.PageSize, cb.GetRootPath()) cachengine.UpdateWriteCountMetric(cb.GetRootPath()) err1 = cb.MaybeWriteCompleted(req.BlockLen) if err1 != nil { return } - totalWritten += int64(readSize) - if totalWritten == req.BlockLen { - if log.EnableDebug() { - log.LogDebugf("action[opCachePutBlock] total write %v", totalWritten) - } + }) + if err == nil { + <-missTaskDone + } + if err1 != nil || err != nil { + break + } else { + if err1 = p.WriteToConn(conn); err1 != nil { + log.LogErrorf("action[opCachePutBlock] reply to conn %v for write data", err1) break } } + totalWritten += int64(readSize) + if totalWritten == req.BlockLen { + if log.EnableDebug() { + log.LogDebugf("action[opCachePutBlock] total write %v", totalWritten) + } + break + } } - }) - if err == nil { - <-missTaskDone } if err1 != nil { f.cacheEngine.DeleteCacheBlock(blockKey) diff --git a/remotecache/flashnode/flashnode_op_test.go b/remotecache/flashnode/flashnode_op_test.go index 4cc9aaf4b..cfeb40c01 100644 --- a/remotecache/flashnode/flashnode_op_test.go +++ b/remotecache/flashnode/flashnode_op_test.go @@ -197,13 +197,10 @@ func testTCPCachePutBlock(t *testing.T) { require.NoError(t, p.WriteToConn(conn)) require.NoError(t, r.ReadFromConn(conn, 3)) require.Equal(t, proto.OpOk, r.ResultCode) - buf := make([]byte, proto.PageSize) + buf := make([]byte, proto.PageSize+4) buf[0] = '{' - _, err := conn.Write(buf[:proto.PageSize]) - require.NoError(t, err) - crcBuf := make([]byte, 4) - binary.BigEndian.PutUint32(crcBuf, crc32.ChecksumIEEE(buf[:proto.PageSize])) - _, err = conn.Write(crcBuf[:4]) + binary.BigEndian.PutUint32(buf[proto.PageSize:], crc32.ChecksumIEEE(buf[:proto.PageSize])) + _, err := conn.Write(buf[:proto.PageSize+4]) require.NoError(t, err) } diff --git a/sdk/remotecache/client.go b/sdk/remotecache/client.go index 3124e5daf..3dcbb6304 100755 --- a/sdk/remotecache/client.go +++ b/sdk/remotecache/client.go @@ -52,14 +52,15 @@ type AddressPingStats struct { } type RemoteCacheClient struct { - flashGroups *btree.BTree - mc *master.MasterClient - conns *util.ConnectPool - hostLatency sync.Map - TTL int64 - ReadTimeout int64 // ms - stopC chan struct{} - wg sync.WaitGroup + flashGroups *btree.BTree + mc *master.MasterClient + conns *util.ConnectPool + hostLatency sync.Map + TTL int64 + ReadTimeout int64 // ms + WriteTimeout int64 + stopC chan struct{} + wg sync.WaitGroup blockSize uint64 clusterEnabled bool @@ -101,6 +102,7 @@ func NewRemoteCacheClient(masters []string, blockSize uint64) (rc *RemoteCacheCl rc.stopC = make(chan struct{}) rc.flashGroups = btree.New(32) rc.ReadTimeout = proto.DefaultRemoteCacheClientReadTimeout + rc.WriteTimeout = proto.DefaultRemoteCacheExtentReadTimeout rc.mc = master.NewMasterClient(masters, false) rc.conns = util.NewConnectPoolWithTimeoutAndCap(5, 500, ConnIdelTimeout, 1) @@ -430,7 +432,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re addr := fg.getFlashHost() if addr == "" { err = fmt.Errorf("getFlashHost failed: can not find host") - log.LogWarnf("FlashGroup reqId(%v) put failed: err(%v)", reqId, err) + log.LogWarnf("FlashGroup fg(%v) reqId(%v) put failed: err(%v)", fg, reqId, err) return } var conn *net.TCPConn @@ -441,7 +443,8 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re } bgTime := stat.BeginStat() defer func() { - rc.conns.PutConnect(conn, err != nil) + forceClose := err != nil && !proto.IsFlashNodeLimitError(err) + rc.conns.PutConnect(conn, forceClose) parts := strings.Split(addr, ":") if len(parts) > 0 && addr != "" { stat.EndStat(fmt.Sprintf("flashPutBlock:%v", parts[0]), err, bgTime, 1) @@ -461,23 +464,30 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re return } replyPacket := NewFlashCacheReply() - if err = replyPacket.ReadFromConnExt(conn, int(rc.ReadTimeout)); err != nil { + if err = replyPacket.ReadFromConnExt(conn, int(rc.WriteTimeout)); err != nil { log.LogWarnf("FlashGroup put: reqId(%v) failed to ReadFromConn, replyPacket(%v), fg host(%v) , err(%v)", reqId, replyPacket, addr, err) return } + if replyPacket.ResultCode != proto.OpOk { + err = fmt.Errorf("%v", string(replyPacket.Data)) + if !proto.IsFlashNodeLimitError(err) { + log.LogWarnf("getPutBlockReply: ResultCode NOK, req(%v) reply(%v) ResultCode(%v)", reqPacket, replyPacket, replyPacket.ResultCode) + } + return + } defer func() { if err != nil { - rc.deleteRemoteBlock(key, conn) + log.LogWarnf("FlashGroup put: reqId(%v) remove key %v by err %v", reqId, key, err) + rc.deleteRemoteBlock(key, addr) } }() var ( totalWritten int64 n int ) - buf := bytespool.Alloc(proto.PageSize) + writeLen := proto.PageSize + 4 + buf := bytespool.Alloc(writeLen) defer bytespool.Free(buf) - crcBuf := bytespool.Alloc(4) - defer bytespool.Free(crcBuf) for { if ctx.Err() != nil { return ctx.Err() @@ -488,18 +498,29 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re } n, err = r.Read(buf[:readSize]) if n != readSize { + log.LogWarnf("FlashGroup put: expected to read %d bytes, but only read %d", readSize, n) return fmt.Errorf("expected to read %d bytes, but only read %d", readSize, n) } if n > 0 { + if log.EnableDebug() { + log.LogDebugf("FlashGroup put: write %d bytes total(%d) to fg", n, totalWritten) + } + binary.BigEndian.PutUint32(buf[proto.PageSize:], crc32.ChecksumIEEE(buf[:proto.PageSize])) conn.SetWriteDeadline(time.Now().Add(proto.WriteDeadlineTime * time.Second)) - if _, err = conn.Write(buf[:proto.PageSize]); err != nil { - log.LogErrorf("wirte data to flashnode get err %v", err) + if _, err = conn.Write(buf[:writeLen]); err != nil { + log.LogErrorf("wirte data and crc to flashnode get err %v", err) return err } - binary.BigEndian.PutUint32(crcBuf, crc32.ChecksumIEEE(buf[:proto.PageSize])) - if _, err = conn.Write(crcBuf[:4]); err != nil { - log.LogErrorf("wirte crc to flashnode get err %v", err) - return err + if err = replyPacket.ReadFromConnExt(conn, int(rc.WriteTimeout)); err != nil { + log.LogWarnf("FlashGroup put data: reqId(%v) failed to ReadFromConn, replyPacket(%v), fg host(%v) , err(%v)", reqId, replyPacket, addr, err) + return + } + if replyPacket.ResultCode != proto.OpOk { + err = fmt.Errorf("%v", string(replyPacket.Data)) + if !proto.IsFlashNodeLimitError(err) { + log.LogWarnf("getPutBlockReply: put data ResultCode NOK, req(%v) reply(%v) ResultCode(%v)", reqPacket, replyPacket, replyPacket.ResultCode) + } + return } totalWritten += int64(n) } @@ -509,6 +530,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re } else if err != nil { return err } else if totalWritten == length { + log.LogDebugf(" total written %v write size %v to fg", totalWritten, readSize) break } } @@ -518,13 +540,29 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re return nil } -func (rc *RemoteCacheClient) deleteRemoteBlock(key string, conn *net.TCPConn) { +func (rc *RemoteCacheClient) deleteRemoteBlock(key string, addr string) { + var ( + err error + conn *net.TCPConn + ) + if conn, err = rc.conns.GetConnect(addr); err != nil { + log.LogWarnf("FlashGroup delete: get connection to curr addr failed, addr(%v) key(%v) err(%v)", addr, key, err) + return + } + defer func() { + rc.conns.PutConnect(conn, err != nil) + }() p := proto.NewPacketReqID() p.Data = ([]byte)(key) p.Opcode = proto.OpFlashNodeCacheDelete p.Size = uint32(len(p.Data)) - if err := p.WriteToConn(conn); err != nil { - log.LogWarnf("FlashGroup put: failed to write to addr(%v) err(%v)", conn.RemoteAddr().String(), err) + if err = p.WriteToConn(conn); err != nil { + log.LogWarnf("FlashGroup delete: failed to write to addr(%v) err(%v)", conn.RemoteAddr().String(), err) + return + } + replyPacket := NewFlashCacheReply() + if err = replyPacket.ReadFromConn(conn, proto.ReadDeadlineTime); err != nil { + log.LogWarnf("FlashGroup delete: failed to ReadFromConn, replyPacket(%v), fg host(%v) err(%v)", replyPacket, addr, err) return } } diff --git a/util/unit.go b/util/unit.go index 26b7ca502..dd93499f9 100644 --- a/util/unit.go +++ b/util/unit.go @@ -22,7 +22,6 @@ import ( "regexp" "strings" - "github.com/cubefs/cubefs/depends/tiglabs/raft/util" "github.com/cubefs/cubefs/util/log" ) @@ -41,7 +40,7 @@ const ( BlockCount = 1024 BlockSize = 65536 * 2 ReadBlockSize = BlockSize - RepairReadBlockSize = 512 * util.KB + RepairReadBlockSize = 512 * KB CacheReadBlockSize = 4 * MB PerBlockCrcSize = 4 ExtentSize = BlockCount * BlockSize @@ -57,7 +56,7 @@ const ( ) const ( - PageSize = 4 * util.KB + PageSize = 4 * KB FallocFLKeepSize = 1 FallocFLPunchHole = 2 )