mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-29 11:15:34 +00:00
* 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>
373 lines
11 KiB
Go
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
|
|
}
|