mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 10:06:14 +00:00
fix(flash): set block allocate size equal to request length
with: #1000292905 Signed-off-by: clinx <chenlin1@oppo.com>
This commit is contained in:
parent
a6fceb08d6
commit
bf3370cc7b
@ -1266,7 +1266,7 @@ func (p *Packet) WriteToConn(c net.Conn) (err error) {
|
||||
return
|
||||
}
|
||||
|
||||
func (p *Packet) WriteToConnForOCS(c net.Conn) (err error) {
|
||||
func (p *Packet) WriteToConnForOCS(c net.Conn, readDiskSize uint32) (err error) {
|
||||
headSize := p.CalcPacketHeaderSize()
|
||||
header, err := Buffers.Get(headSize)
|
||||
if err != nil {
|
||||
@ -1292,7 +1292,7 @@ func (p *Packet) WriteToConnForOCS(c net.Conn) (err error) {
|
||||
}
|
||||
if _, err = c.Write(p.Arg[:int(p.ArgLen)]); err == nil {
|
||||
if p.Data != nil && p.Size != 0 {
|
||||
_, err = c.Write(p.Data[:])
|
||||
_, err = c.Write(p.Data[:readDiskSize])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@ -211,7 +211,7 @@ func (cb *CacheBlock) Read(ctx context.Context, data []byte, offset, size int64,
|
||||
}
|
||||
|
||||
if cb.sourceType == SourceTypeBlock {
|
||||
realSize = proto.CACHE_BLOCK_PACKET_SIZE
|
||||
realSize = size
|
||||
}
|
||||
|
||||
log.LogDebugf("action[Read] read cache block:%v, offset:%d, allocSize:%d, usedSize:%d", cb.blockKey, offset, cb.allocSize, cb.usedSize)
|
||||
@ -923,9 +923,6 @@ func CalcAllocSizeV2(reqLen int) (int, error) {
|
||||
if reqLen > proto.CACHE_OBJECT_BLOCK_SIZE {
|
||||
return 0, fmt.Errorf("invalid block size: %d", reqLen)
|
||||
}
|
||||
if reqLen%proto.CACHE_BLOCK_PACKET_SIZE != 0 {
|
||||
reqLen = (reqLen/proto.CACHE_BLOCK_PACKET_SIZE + 1) * proto.CACHE_BLOCK_PACKET_SIZE
|
||||
}
|
||||
return reqLen, nil
|
||||
}
|
||||
|
||||
|
||||
@ -389,32 +389,32 @@ func (f *FlashNode) opCachePutBlock(conn net.Conn, p *proto.Packet) (err error)
|
||||
}
|
||||
readSize = proto.CACHE_BLOCK_PACKET_SIZE
|
||||
missTaskDone := make(chan struct{})
|
||||
err = f.limitWrite.Run(writeLen, true, func() {
|
||||
if totalWritten+int64(readSize) > req.BlockLen {
|
||||
readSize = int(req.BlockLen - totalWritten)
|
||||
}
|
||||
err = f.limitWrite.Run(readSize, true, func() {
|
||||
defer func() {
|
||||
close(missTaskDone)
|
||||
}()
|
||||
if totalWritten+int64(readSize) > req.BlockLen {
|
||||
readSize = int(req.BlockLen - totalWritten)
|
||||
}
|
||||
if n, err1 = io.ReadFull(conn, buf[:writeLen]); err1 != nil {
|
||||
if n, err1 = io.ReadFull(conn, buf[:readSize+proto.CACHE_BLOCK_CRC_SIZE]); err1 != nil {
|
||||
log.LogWarnf(logPrefix+" blockkey %v read data and crc conn %v", blockKey, err1)
|
||||
return
|
||||
}
|
||||
if n != writeLen {
|
||||
if n != readSize+proto.CACHE_BLOCK_CRC_SIZE {
|
||||
err1 = syscall.EBADMSG
|
||||
return
|
||||
}
|
||||
err1 = cb.WriteAtV2(&proto.FlashWriteParam{
|
||||
Offset: totalWritten,
|
||||
Size: req.BlockLen,
|
||||
Data: buf[:proto.CACHE_BLOCK_PACKET_SIZE],
|
||||
Crc: buf[proto.CACHE_BLOCK_PACKET_SIZE:writeLen],
|
||||
Data: buf[:readSize],
|
||||
Crc: buf[readSize : readSize+proto.CACHE_BLOCK_CRC_SIZE],
|
||||
DataSize: int64(readSize),
|
||||
})
|
||||
if err1 != nil {
|
||||
return
|
||||
}
|
||||
cachengine.UpdateWriteBytesMetric(proto.CACHE_BLOCK_PACKET_SIZE, cb.GetRootPath())
|
||||
cachengine.UpdateWriteBytesMetric(uint64(readSize), cb.GetRootPath())
|
||||
cachengine.UpdateWriteCountMetric(cb.GetRootPath())
|
||||
err1 = cb.MaybeWriteCompleted(req.BlockLen)
|
||||
if err1 != nil {
|
||||
@ -792,38 +792,36 @@ func (f *FlashNode) doObjectReadRequest(ctx context.Context, conn net.Conn, req
|
||||
var errInner error
|
||||
buf := bytespool.Alloc(proto.CACHE_BLOCK_PACKET_SIZE)
|
||||
defer bytespool.Free(buf)
|
||||
var alignedOffset int64
|
||||
var readDiskSize uint32
|
||||
readAndReply := func() {
|
||||
reply := proto.NewPacket()
|
||||
reply.ReqID = p.ReqID
|
||||
reply.StartT = p.StartT
|
||||
reply.Data = buf
|
||||
|
||||
alignedOffset := offset / proto.CACHE_BLOCK_PACKET_SIZE * proto.CACHE_BLOCK_PACKET_SIZE
|
||||
reply.KernelOffset = uint64(offset)
|
||||
reply.ExtentOffset = offset - alignedOffset
|
||||
p.Size = proto.CACHE_BLOCK_PACKET_SIZE
|
||||
p.ExtentOffset = offset
|
||||
|
||||
reply.CRC, errInner = block.Read(ctx, reply.Data[:], alignedOffset, proto.CACHE_BLOCK_PACKET_SIZE, f.waitForCacheBlock, true)
|
||||
reply.CRC, errInner = block.Read(ctx, reply.Data[:readDiskSize], alignedOffset, int64(readDiskSize), f.waitForCacheBlock, true)
|
||||
if errInner != nil {
|
||||
return
|
||||
}
|
||||
p.CRC = reply.CRC
|
||||
realNeedSize := uint32(util.Min(int(proto.CACHE_BLOCK_PACKET_SIZE-reply.ExtentOffset), int(end-offset)))
|
||||
|
||||
reply.Size = realNeedSize
|
||||
reply.ResultCode = proto.OpOk
|
||||
reply.Opcode = p.Opcode
|
||||
p.ResultCode = proto.OpOk
|
||||
|
||||
bgTime := stat.BeginStat()
|
||||
if errInner = reply.WriteToConnForOCS(conn); errInner != nil {
|
||||
if errInner = reply.WriteToConnForOCS(conn, readDiskSize); errInner != nil {
|
||||
log.LogErrorf("%s key:[%s] %s", action, block.GetBlockKey(),
|
||||
reply.LogMessage(reply.GetOpMsg(), conn.RemoteAddr().String(), reply.StartT, errInner))
|
||||
return
|
||||
}
|
||||
stat.EndStat("HitCacheRead:ReplyToClient", errInner, bgTime, 1)
|
||||
offset = alignedOffset + proto.CACHE_BLOCK_PACKET_SIZE
|
||||
offset = alignedOffset + int64(readDiskSize)
|
||||
if log.EnableInfo() {
|
||||
log.LogInfof("%s ReqID[%d] key:[%s] reply[%s] block[%s]", action, p.ReqID, block.GetBlockKey(),
|
||||
reply.LogMessage(reply.GetOpMsg(), conn.RemoteAddr().String(), reply.StartT, errInner), block.GetBlockKey())
|
||||
@ -831,11 +829,13 @@ func (f *FlashNode) doObjectReadRequest(ctx context.Context, conn net.Conn, req
|
||||
}
|
||||
var keepAlive bool
|
||||
for {
|
||||
alignedOffset = offset / proto.CACHE_BLOCK_PACKET_SIZE * proto.CACHE_BLOCK_PACKET_SIZE
|
||||
readDiskSize = uint32(util.Min(proto.CACHE_BLOCK_PACKET_SIZE, int(end-alignedOffset)))
|
||||
if !keepAlive {
|
||||
err = f.limitRead.RunNoWait(proto.CACHE_BLOCK_PACKET_SIZE, false, readAndReply)
|
||||
err = f.limitRead.RunNoWait(int(readDiskSize), false, readAndReply)
|
||||
keepAlive = true
|
||||
} else {
|
||||
err = f.limitRead.Run(proto.CACHE_BLOCK_PACKET_SIZE, true, readAndReply)
|
||||
err = f.limitRead.Run(int(readDiskSize), true, readAndReply)
|
||||
}
|
||||
if err != nil {
|
||||
return
|
||||
|
||||
@ -187,9 +187,10 @@ func testTCPCachePutBlock(t *testing.T) {
|
||||
require.NoError(t, r.ReadFromConn(conn, 3))
|
||||
require.Equal(t, proto.OpErr, r.ResultCode) // lack parameter
|
||||
key := "testTCPCacheWrite_key"
|
||||
bLen := 1
|
||||
req := &proto.PutBlockHead{
|
||||
UniKey: key,
|
||||
BlockLen: 1,
|
||||
BlockLen: int64(bLen),
|
||||
TTL: 0,
|
||||
}
|
||||
_ = p.MarshalDataPb(req)
|
||||
@ -199,8 +200,8 @@ func testTCPCachePutBlock(t *testing.T) {
|
||||
require.Equal(t, proto.OpOk, r.ResultCode)
|
||||
buf := make([]byte, proto.CACHE_BLOCK_PACKET_SIZE+4)
|
||||
buf[0] = '{'
|
||||
binary.BigEndian.PutUint32(buf[proto.CACHE_BLOCK_PACKET_SIZE:], crc32.ChecksumIEEE(buf[:proto.CACHE_BLOCK_PACKET_SIZE]))
|
||||
_, err := conn.Write(buf[:proto.CACHE_BLOCK_PACKET_SIZE+4])
|
||||
binary.BigEndian.PutUint32(buf[bLen:], crc32.ChecksumIEEE(buf[:bLen]))
|
||||
_, err := conn.Write(buf[:bLen+4])
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
|
||||
@ -670,9 +670,9 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
|
||||
if log.EnableDebug() {
|
||||
log.LogDebugf("FlashGroup put: write %d bytes total(%d) to fg", bufferOffset, totalWritten)
|
||||
}
|
||||
binary.BigEndian.PutUint32(buf[proto.CACHE_BLOCK_PACKET_SIZE:], crc32.ChecksumIEEE(buf[:proto.CACHE_BLOCK_PACKET_SIZE]))
|
||||
binary.BigEndian.PutUint32(buf[bufferOffset:], crc32.ChecksumIEEE(buf[:bufferOffset]))
|
||||
conn.SetWriteDeadline(time.Now().Add(proto.WriteDeadlineTime * time.Second))
|
||||
if _, err = conn.Write(buf[:writeLen]); err != nil {
|
||||
if _, err = conn.Write(buf[:bufferOffset+proto.CACHE_BLOCK_CRC_SIZE]); err != nil {
|
||||
log.LogErrorf("wirte data and crc to flashnode get err %v", err)
|
||||
return fmt.Errorf(proto.ErrorWriteDataAndCRCToFlashNodeTpl, err)
|
||||
}
|
||||
@ -713,7 +713,7 @@ func (rc *RemoteCacheClient) processPutReply(conn *net.TCPConn, ch *proto.CoonHa
|
||||
var err error
|
||||
replyPacket := NewFlashCacheReply()
|
||||
for range ch.WaitAckChan {
|
||||
if err = replyPacket.ReadFromConnExt(conn, int(rc.WriteTimeout)); err != nil {
|
||||
if err = replyPacket.ReadFromConn(conn, proto.ReadDeadlineTime); err != nil {
|
||||
log.LogWarnf("FlashGroup put data: reqId(%v) failed to ReadFromConn, replyPacket(%v) , err(%v)", reqId, replyPacket, err)
|
||||
ch.RemoteError = err
|
||||
return
|
||||
@ -1169,14 +1169,16 @@ func (reader *RemoteCacheReader) read(p []byte) (n int, err error) {
|
||||
}
|
||||
return
|
||||
}
|
||||
_, err = io.ReadFull(reader.conn, p)
|
||||
alignedOffset := reader.currOffset / proto.CACHE_BLOCK_PACKET_SIZE * proto.CACHE_BLOCK_PACKET_SIZE
|
||||
readPackageSize := uint32(util.Min(proto.CACHE_BLOCK_PACKET_SIZE, int(reader.endOffset-alignedOffset)))
|
||||
_, err = io.ReadFull(reader.conn, p[:readPackageSize])
|
||||
if err != nil {
|
||||
log.LogErrorf("RemoteCacheReader:reqID(%v) Read err(%v)", reader.reqID, err)
|
||||
return
|
||||
}
|
||||
|
||||
// check crc
|
||||
actualCrc := crc32.ChecksumIEEE(p)
|
||||
actualCrc := crc32.ChecksumIEEE(p[:readPackageSize])
|
||||
if actualCrc != reply.CRC {
|
||||
err = fmt.Errorf(proto.ErrorInconsistentCRCObjectTpl, reply.KernelOffset, reply.ExtentOffset, reply.CRC, actualCrc)
|
||||
log.LogErrorf("RemoteCacheReader:Read check crc failed reqID(%v) offset(%v) extentOffset(%v) expect(%v) actualCrc(%v)",
|
||||
@ -1188,7 +1190,7 @@ func (reader *RemoteCacheReader) read(p []byte) (n int, err error) {
|
||||
if log.EnableDebug() {
|
||||
log.LogDebugf("RemoteCacheReader:Read reqID(%v) offset(%v) extentOffset(%v) expectLen(%v) p.len(%v) cost(%v)", reader.reqID, reply.KernelOffset, reply.ExtentOffset, expectLen, len(p), time.Since(*bg))
|
||||
}
|
||||
atomic.StoreInt64(&reader.alreadyReadLen, reader.alreadyReadLen+expectLen)
|
||||
atomic.AddInt64(&reader.alreadyReadLen, expectLen)
|
||||
return int(expectLen), nil
|
||||
}
|
||||
|
||||
@ -1206,6 +1208,9 @@ func (reader *RemoteCacheReader) Read(p []byte) (n int, err error) {
|
||||
if len(p) == 0 {
|
||||
return 0, nil
|
||||
}
|
||||
if log.EnableDebug() {
|
||||
log.LogDebugf("RemoteCacheReader:reqID(%v) start(%v) readsize(%v)", reader.reqID, reader.currOffset, len(p))
|
||||
}
|
||||
if reader.closed {
|
||||
log.LogErrorf("read from close reader reqID(%v)", reader.reqID)
|
||||
return 0, fmt.Errorf(proto.ErrorReadFromCloseReaderTpl, reader.reqID)
|
||||
|
||||
Loading…
Reference in New Issue
Block a user