Files
seaweedfs/weed/storage/volume_read.go
T
yi111GitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Chris LuDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
5389f61cef volume server: do not finish a GET when the needle CRC mismatches (#11464)
* volume server: do not finish a GET when the needle CRC mismatches

A streamed full-needle read compared the CRC only after every page had been written. Once the response buffer flushed, the client already had a completed 200 and the corrupt bytes. Hold the last page until the checksum matches, and if an earlier page has already been flushed, abort the connection instead of calling http.Error.

Fixes #11459

* volume server: abort partial-content bodies on write error too

The non-Range path drops the unflushed tail and aborts on a mid-body
error; the single-range and multi-range paths still flushed it after
WriteHeader(206) was committed, delivering corrupt bytes as a complete
body.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* volume server: assert the started 200 is aborted in the write-error test

The test previously returned on any request error, so it passed without
verifying the abort. It now asserts the client got the committed 200
headers and then a failed body read. Also trims comments.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-09-26 16:00:22 +08:00

373 lines
11 KiB
Go

package storage
import (
"fmt"
"io"
"time"
"github.com/seaweedfs/seaweedfs/weed/util/mem"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/stats"
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
. "github.com/seaweedfs/seaweedfs/weed/storage/types"
)
const PagedReadLimit = 1024 * 1024
// read fills in Needle content by looking up n.Id from NeedleMapper
func (v *Volume) readNeedle(n *needle.Needle, readOption *ReadOption, onReadSizeFn func(size Size)) (count int, err error) {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return 0, err
}
if v.nm == nil {
glog.V(0).Infof("volume %d: needle map not loaded; read returns not-found", v.Id)
return -1, ErrorNotFound
}
nv, ok := v.nm.Get(n.Id)
if !ok || nv.Offset.IsZero() {
return -1, ErrorNotFound
}
readSize := nv.Size
if readSize.IsDeleted() {
if readOption != nil && readOption.ReadDeleted && readSize != TombstoneFileSize {
glog.V(3).Infof("reading deleted %s", n.String())
stats.VolumeServerHandlerCounter.WithLabelValues(stats.ReadDeletedNeedle).Inc()
readSize = -readSize
} else {
return -1, ErrorDeleted
}
}
if readSize == 0 {
return 0, nil
}
if onReadSizeFn != nil {
onReadSizeFn(readSize)
}
if readOption != nil && readOption.AttemptMetaOnly && readSize > PagedReadLimit {
readOption.VolumeRevision = v.SuperBlock.CompactionRevision
err = n.ReadNeedleMeta(v.DataBackend, nv.Offset.ToActualOffset(), readSize, v.Version())
if err == needle.ErrorSizeMismatch && OffsetSize == 4 {
readOption.IsOutOfRange = true
err = n.ReadNeedleMeta(v.DataBackend, nv.Offset.ToActualOffset()+int64(MaxPossibleVolumeSize), readSize, v.Version())
}
if err != nil {
return 0, err
}
if !n.IsCompressed() && !n.IsChunkedManifest() {
readOption.IsMetaOnly = true
}
}
if readOption == nil || !readOption.IsMetaOnly {
err = n.ReadData(v.DataBackend, nv.Offset.ToActualOffset(), readSize, v.Version())
v.checkReadWriteError(err)
if err != nil {
return 0, err
}
}
count = int(n.DataSize)
if !n.HasTtl() {
return
}
ttlMinutes := n.Ttl.Minutes()
if ttlMinutes == 0 {
return
}
if !n.HasLastModifiedDate() {
return
}
if time.Now().Before(time.Unix(0, int64(n.AppendAtNs)).Add(time.Duration(ttlMinutes) * time.Minute)) {
return
}
return -1, ErrorNotFound
}
// read needle at a specific offset
func (v *Volume) readNeedleMetaAt(n *needle.Needle, offset int64, size int32) (err error) {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return err
}
// read deleted needle meta data
if size < 0 {
size = 0
}
err = n.ReadNeedleMeta(v.DataBackend, offset, Size(size), v.Version())
if err == needle.ErrorSizeMismatch && OffsetSize == 4 {
err = n.ReadNeedleMeta(v.DataBackend, offset+int64(MaxPossibleVolumeSize), Size(size), v.Version())
}
if err != nil {
return err
}
return nil
}
// read fills in Needle content by looking up n.Id from NeedleMapper
func (v *Volume) readNeedleDataInto(n *needle.Needle, readOption *ReadOption, writer io.Writer, offset int64, size int64) (err error) {
if !readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
}
if readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
}
if err := v.UnavailableError(); err != nil {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
return err
}
if v.nm == nil {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
glog.V(0).Infof("volume %d: needle map not loaded; read returns not-found", v.Id)
return ErrorNotFound
}
nv, ok := v.nm.Get(n.Id)
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
if !ok || nv.Offset.IsZero() {
return ErrorNotFound
}
readSize := nv.Size
if readSize.IsDeleted() {
if readOption != nil && readOption.ReadDeleted && readSize != TombstoneFileSize {
glog.V(3).Infof("reading deleted %s", n.String())
readSize = -readSize
} else {
return ErrorDeleted
}
}
if readSize == 0 {
return nil
}
actualOffset := nv.Offset.ToActualOffset()
if readOption.IsOutOfRange {
actualOffset += int64(MaxPossibleVolumeSize)
}
buf := mem.Allocate(min(readOption.ReadBufferSize, int(size)))
// A full-needle read holds the last page back in `pending` until the CRC
// matches; `spare` is a second buffer swapped in so the held page is not
// overwritten while streaming.
var spare []byte
defer func() {
mem.Free(buf)
if spare != nil {
mem.Free(spare)
}
}()
// read needle data
crc := needle.CRC(0)
var pending []byte
checkCRC := offset == 0 && size == int64(n.DataSize)
for x := offset; x < offset+size; x += int64(len(buf)) {
if readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
if err := v.UnavailableError(); err != nil {
v.dataFileAccessLock.RUnlock()
return err
}
}
// possibly re-read needle offset if volume is compacted
if readOption.VolumeRevision != v.SuperBlock.CompactionRevision {
if v.nm == nil {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
glog.V(0).Infof("volume %d: needle map not loaded mid-read", v.Id)
return ErrorNotFound
}
// the volume is compacted
nv, ok = v.nm.Get(n.Id)
if !ok || nv.Offset.IsZero() {
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
return ErrorNotFound
}
actualOffset = nv.Offset.ToActualOffset()
readOption.VolumeRevision = v.SuperBlock.CompactionRevision
}
count, err := n.ReadNeedleData(v.DataBackend, actualOffset, buf, x)
if readOption.HasSlowRead {
v.dataFileAccessLock.RUnlock()
}
// Thread the underlying read error through the EIO tracker.
// Without this, large/range GETs through readNeedleDataInto
// would never trip IoErrorTolerance even on a failing disk.
// io.EOF is treated as a clean end-of-stream below, not an
// error.
if err != nil && err != io.EOF {
v.checkReadWriteError(err)
}
toWrite := min(count, int(offset+size-x))
if toWrite > 0 {
crc = crc.Update(buf[0:toWrite])
// The CRC is known only after the last byte; hold each page until
// the next one is read so a bad needle is never fully written.
if checkCRC {
if pending != nil {
if _, err = writer.Write(pending); err != nil {
return fmt.Errorf("ReadNeedleData write: %w", err)
}
}
pending = buf[:toWrite]
if x+int64(len(buf)) < offset+size {
if spare == nil {
spare = mem.Allocate(len(buf))
}
buf, spare = spare, buf
}
} else if _, err = writer.Write(buf[0:toWrite]); err != nil {
return fmt.Errorf("ReadNeedleData write: %w", err)
}
}
if err != nil {
if err == io.EOF {
err = nil
break
}
return fmt.Errorf("ReadNeedleData: %w", err)
}
if count <= 0 {
break
}
}
// Whole-needle read completed without a backend error — clear any
// pending EIO streak. If a non-EIO failure happens later (CRC etc.)
// we still return that error to the caller, but the disk itself
// produced clean bytes.
v.checkReadWriteError(nil)
if checkCRC && (n.Checksum != crc && uint32(n.Checksum) != crc.Value()) {
// the crc.Value() function is to be deprecated. this double checking is for backward compatibility
// with seaweed version using crc.Value() instead of uint32(crc), which appears in commit 056c480eb
// and switch appeared in version 3.09.
// pending is the last page and is intentionally not written.
stats.VolumeServerHandlerCounter.WithLabelValues(stats.ErrorCRC).Inc()
return fmt.Errorf("ReadNeedleData checksum %v expected %v for Needle: %v,%v", crc, n.Checksum, v.Id, n)
}
if pending != nil {
if _, err = writer.Write(pending); err != nil {
return fmt.Errorf("ReadNeedleData write: %w", err)
}
}
return nil
}
func min(x, y int) int {
if x < y {
return x
}
return y
}
// read fills in Needle content by looking up n.Id from NeedleMapper
func (v *Volume) ReadNeedleBlob(offset int64, size Size) ([]byte, error) {
// A deletion marker is not a record length; reject it before taking the lock.
if size.IsDeleted() {
return nil, fmt.Errorf("invalid needle size %d", size)
}
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
if err := v.UnavailableError(); err != nil {
return nil, err
}
blob, err := needle.ReadNeedleBlob(v.DataBackend, offset, size, v.Version())
v.checkReadWriteError(err)
return blob, err
}
type VolumeFileScanner interface {
VisitSuperBlock(super_block.SuperBlock) error
ReadNeedleBody() bool
VisitNeedle(n *needle.Needle, offset int64, needleHeader, needleBody []byte) error
}
func ScanVolumeFile(dirname string, collection string, id needle.VolumeId,
needleMapKind NeedleMapKind,
volumeFileScanner VolumeFileScanner) (err error) {
var v *Volume
if v, err = loadVolumeWithoutIndex(dirname, collection, id, needleMapKind, needle.GetCurrentVersion()); err != nil {
return fmt.Errorf("failed to load volume %d: %w", id, err)
}
if err = volumeFileScanner.VisitSuperBlock(v.SuperBlock); err != nil {
return fmt.Errorf("failed to process volume %d super block: %w", id, err)
}
defer v.Close()
version := v.Version()
offset := int64(v.SuperBlock.BlockSize())
return ScanVolumeFileFrom(version, v.DataBackend, offset, volumeFileScanner)
}
func ScanVolumeFileFrom(version needle.Version, datBackend backend.BackendStorageFile, offset int64, volumeFileScanner VolumeFileScanner) (err error) {
n, nh, rest, e := needle.ReadNeedleHeader(datBackend, version, offset)
if e != nil {
if e == io.EOF {
return nil
}
return fmt.Errorf("cannot read %s at offset %d: %w", datBackend.Name(), offset, e)
}
for n != nil {
var needleBody []byte
if volumeFileScanner.ReadNeedleBody() {
// println("needle", n.Id.String(), "offset", offset, "size", n.Size, "rest", rest)
if needleBody, err = n.ReadNeedleBody(datBackend, version, offset+NeedleHeaderSize, rest); err != nil {
glog.V(0).Infof("cannot read needle head [%d, %d) body [%d, %d) body length %d: %v", offset, offset+NeedleHeaderSize, offset+NeedleHeaderSize, offset+NeedleHeaderSize+rest, rest, err)
// err = fmt.Errorf("cannot read needle body: %v", err)
// return
}
}
err := volumeFileScanner.VisitNeedle(n, offset, nh, needleBody)
if err == io.EOF {
return nil
}
if err != nil {
glog.V(0).Infof("visit needle error: %v", err)
return fmt.Errorf("visit needle error: %w", err)
}
// A corrupt header can carry a size so negative that the record length
// is zero or less; the scan cannot advance past it.
recordSize := NeedleHeaderSize + rest
if recordSize <= 0 {
return fmt.Errorf("%s: needle header at offset %d has size %d, record length %d: %w", datBackend.Name(), offset, n.Size, recordSize, needle.ErrorCorrupted)
}
offset += recordSize
glog.V(4).Infof("==> new entry offset %d", offset)
if n, nh, rest, err = needle.ReadNeedleHeader(datBackend, version, offset); err != nil {
if err == io.EOF {
return nil
}
return fmt.Errorf("cannot read needle header at offset %d: %w", offset, err)
}
glog.V(4).Infof("new entry needle size:%d rest:%d", n.Size, rest)
}
return nil
}