From f2cfda41261a1d7917837c8cf22198f0d1e187d3 Mon Sep 17 00:00:00 2001 From: slasher Date: Thu, 2 Jan 2025 16:29:42 +0800 Subject: [PATCH] feat(common): crc32 decoder with load mode return head and data . #22922161 + #3718 Signed-off-by: slasher --- blobstore/common/crc32block/sized_coder.go | 81 +++++++++++++++- .../common/crc32block/sized_coder_test.go | 94 +++++++++++++++++++ 2 files changed, 170 insertions(+), 5 deletions(-) diff --git a/blobstore/common/crc32block/sized_coder.go b/blobstore/common/crc32block/sized_coder.go index 404e6f940..e99c838b3 100644 --- a/blobstore/common/crc32block/sized_coder.go +++ b/blobstore/common/crc32block/sized_coder.go @@ -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)) } } diff --git a/blobstore/common/crc32block/sized_coder_test.go b/blobstore/common/crc32block/sized_coder_test.go index 6405d71f1..77455209b 100644 --- a/blobstore/common/crc32block/sized_coder_test.go +++ b/blobstore/common/crc32block/sized_coder_test.go @@ -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 }