mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(common): crc32 decoder with load mode
return head and data . #22922161 + #3718 Signed-off-by: slasher <shenjie1@oppo.com>
This commit is contained in:
parent
427d0ae09c
commit
f2cfda4126
@ -16,6 +16,7 @@ package crc32block
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"fmt"
|
||||
"hash"
|
||||
"hash/crc32"
|
||||
"io"
|
||||
@ -24,9 +25,10 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
ModeEncode uint8 = 1
|
||||
ModeCheck uint8 = 2
|
||||
ModeDecode uint8 = 3
|
||||
ModeEncode uint8 = 1 // encode mode
|
||||
ModeCheck uint8 = 2 // check, return buffer with head, crc and tail
|
||||
ModeDecode uint8 = 3 // decode, return data buffer
|
||||
ModeLoad uint8 = 4 // load, return buffer with head, should with section
|
||||
|
||||
_alignment = transport.Alignment
|
||||
_alignmentMask = _alignment - 1
|
||||
@ -264,6 +266,73 @@ func (r *sizedCoder) decodeRead(p []byte) (nn int, err error) {
|
||||
return
|
||||
}
|
||||
|
||||
// decodeLoad writes origin head and data to the parameter p.
|
||||
func (r *sizedCoder) decodeLoad(p []byte) (nn int, err error) {
|
||||
if r.err != nil {
|
||||
return 0, r.err
|
||||
}
|
||||
if r.remain <= 0 {
|
||||
return 0, io.EOF
|
||||
}
|
||||
if !r.section {
|
||||
return 0, fmt.Errorf("crc32block: should load with sectioned")
|
||||
}
|
||||
if len(p) == 0 || len(p)%_alignment != 0 {
|
||||
return 0, fmt.Errorf("crc32block: should aligned buffer, but %d", len(p))
|
||||
}
|
||||
|
||||
tryRead := r.payload + crc32.Size - r.nx
|
||||
extra := tryRead - r.padhead - crc32.Size - r.padtail - r.remain
|
||||
if extra >= 0 {
|
||||
tryRead -= extra
|
||||
}
|
||||
if len(p) < tryRead {
|
||||
return 0, fmt.Errorf("crc32block: should enough buffer %d, but %d", tryRead, len(p))
|
||||
}
|
||||
|
||||
var n int
|
||||
n, err = r.ReadCloser.Read(p[:tryRead])
|
||||
if n != tryRead {
|
||||
return 0, fmt.Errorf("crc32block: short of read should %d, buf %d", tryRead, n)
|
||||
}
|
||||
if r.padhead > 0 {
|
||||
p = p[r.padhead:]
|
||||
nn += r.padhead // notice to caller, the buffer has head
|
||||
r.nx += r.padhead
|
||||
n -= r.padhead
|
||||
r.padhead = 0
|
||||
}
|
||||
|
||||
n -= crc32.Size
|
||||
if extra >= 0 { // last block
|
||||
n -= r.padtail
|
||||
r.padtail = 0
|
||||
}
|
||||
copy(r.cell[:], p[n:])
|
||||
|
||||
nn += n
|
||||
r.crc32.Write(p[:n])
|
||||
r.nx += n
|
||||
r.remain -= n
|
||||
|
||||
if r.nx == r.payload || r.remain == 0 {
|
||||
act := r.crc32.Sum(nil)
|
||||
if !bytes.Equal(r.cell[:], act) {
|
||||
r.err = ErrMismatchedCrc
|
||||
return 0, r.err
|
||||
}
|
||||
|
||||
r.crc32.Reset()
|
||||
r.nx = 0
|
||||
r.cx = -1
|
||||
|
||||
if r.section && r.remain > 0 { // return sectioned error
|
||||
err = transport.ErrFrameContinue
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (r *sizedCoder) Read(p []byte) (nn int, err error) {
|
||||
var n int
|
||||
for len(p) > 0 {
|
||||
@ -272,10 +341,12 @@ func (r *sizedCoder) Read(p []byte) (nn int, err error) {
|
||||
n, err = r.encodeRead(p)
|
||||
case ModeDecode:
|
||||
n, err = r.decodeRead(p)
|
||||
case ModeLoad:
|
||||
return r.decodeLoad(p)
|
||||
case ModeCheck:
|
||||
panic("crc32block: implement checker with WriterTo")
|
||||
default:
|
||||
panic("crc32block: unknow mode with Reader")
|
||||
panic(fmt.Sprintf("crc32block: unknow mode %d with Reader", r.mode))
|
||||
}
|
||||
nn += n
|
||||
p = p[n:]
|
||||
@ -305,7 +376,7 @@ func (r *sizedCoderWriter) Write(p []byte) (nn int, err error) {
|
||||
case ModeCheck:
|
||||
return r.checkWrite(p)
|
||||
default:
|
||||
panic("crc32block: unknow mode with WriterTo")
|
||||
panic(fmt.Sprintf("crc32block: unknow mode %d with WriterTo", r.mode))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -393,6 +393,100 @@ func TestSizedCoderChecker(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSizedCoderLoader(t *testing.T) {
|
||||
{
|
||||
rc := NewSizedCoder(nil, 0, 0, 32<<10, ModeLoad, true)
|
||||
_, err := rc.Read(make([]byte, 1))
|
||||
require.ErrorIs(t, io.EOF, err)
|
||||
}
|
||||
{
|
||||
rc := NewSizedCoder(nil, 1, 0, 32<<10, ModeLoad, false)
|
||||
_, err := rc.Read(make([]byte, 1))
|
||||
require.Error(t, err)
|
||||
}
|
||||
{
|
||||
rc := NewSizedCoder(io.NopCloser(bytes.NewReader(make([]byte, 10))),
|
||||
1, 0, 32<<10, ModeLoad, true)
|
||||
_, err := rc.Read(make([]byte, 1))
|
||||
require.Error(t, err)
|
||||
_, err = rc.Read(make([]byte, _alignment))
|
||||
require.Error(t, err)
|
||||
}
|
||||
{
|
||||
size := int64(_alignmentMask)
|
||||
buf := make([]byte, size)
|
||||
re := NewSizedCoder(io.NopCloser(bytes.NewReader(buf)),
|
||||
size, 0, 32<<10, ModeEncode, false)
|
||||
ebuf, err := io.ReadAll(re)
|
||||
require.NoError(t, err)
|
||||
ebuf[_alignmentMask+2]++
|
||||
rd := NewSizedCoder(io.NopCloser(bytes.NewReader(ebuf)),
|
||||
size, 0, 32<<10, ModeLoad, true)
|
||||
_, err = rd.Read(make([]byte, _alignment))
|
||||
require.Error(t, err)
|
||||
_, err = rd.Read(make([]byte, _alignment*2))
|
||||
require.ErrorIs(t, err, ErrMismatchedCrc)
|
||||
_, err = rd.Read(make([]byte, _alignment*2))
|
||||
require.ErrorIs(t, err, ErrMismatchedCrc)
|
||||
}
|
||||
for _, size := range []int64{_alignment - 1, _alignment, _alignment + 1} {
|
||||
buf := make([]byte, size)
|
||||
re := NewSizedCoder(io.NopCloser(bytes.NewReader(buf)),
|
||||
size, 0, 32<<10, ModeEncode, false)
|
||||
ebuf, err := io.ReadAll(re)
|
||||
require.NoError(t, err)
|
||||
rd := NewSizedCoder(io.NopCloser(bytes.NewReader(ebuf)),
|
||||
size, 0, 32<<10, ModeLoad, true)
|
||||
n, err := rd.Read(make([]byte, _alignment*2))
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int(size), n)
|
||||
}
|
||||
|
||||
size := int64(1<<20) + 19
|
||||
buff := make([]byte, size)
|
||||
crand.Read(buff)
|
||||
rbuff := make([]byte, gBlockSize)
|
||||
|
||||
run := func(actual, stable int64) {
|
||||
buf := buff[0:actual:actual]
|
||||
|
||||
re := NewSizedCoder(io.NopCloser(bytes.NewReader(buf)),
|
||||
actual, stable, gBlockSize, ModeEncode, false)
|
||||
ebuf, err := io.ReadAll(re)
|
||||
require.NoError(t, err)
|
||||
|
||||
rd := NewSizedCoder(io.NopCloser(bytes.NewReader(ebuf)),
|
||||
actual, stable, gBlockSize, ModeLoad, true)
|
||||
|
||||
var nn int
|
||||
crc := crc32.NewIEEE()
|
||||
head := int((stable % BlockPayload(gBlockSize)) % _alignment)
|
||||
for {
|
||||
n, err := rd.Read(rbuff)
|
||||
crc.Write(rbuff[head:n])
|
||||
nn += n - head
|
||||
head = 0
|
||||
if err == transport.ErrFrameContinue {
|
||||
continue
|
||||
}
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
require.NoError(t, err)
|
||||
}
|
||||
require.Equal(t, len(buf), nn)
|
||||
require.Equal(t, crc32.ChecksumIEEE(buf), crc.Sum32())
|
||||
}
|
||||
|
||||
run(size, 0)
|
||||
run(1, size-1)
|
||||
run(1999, 817374)
|
||||
run(2000, 11223344)
|
||||
for range [100]struct{}{} {
|
||||
run(mrand.Int63n(size), mrand.Int63n(1<<20))
|
||||
}
|
||||
}
|
||||
|
||||
type noneReadWriter struct{}
|
||||
|
||||
func (noneReadWriter) Read(p []byte) (int, error) { return len(p), nil }
|
||||
|
||||
Loading…
Reference in New Issue
Block a user