Files
seaweedfs/weed/storage/volume_read.go
T
Chris LuandGitHub ce1e0dc30a s3api: don't delete chunks when CreateEntry outcome is ambiguous (#11376)
* s3api: map ambiguous filer transport errors to retryable 503

Canceled, DeadlineExceeded and Unavailable can be returned after the
filer applied the write, so the outcome is ambiguous. Reporting them as
a 4xx tells the client not to retry; report ServiceUnavailable instead.

* s3api: verify entry existence before deleting orphaned chunks

A failed CreateEntry can still have landed on the filer when the error
is a transport failure, and entryCreated=false would tombstone chunks a
live entry references, leaving a dangling pointer that survives only
because reads pass readDeleted=true until vacuum reclaims the needle.

Before deleting, look the entry up: if it is stored with the same
chunks, the write succeeded; if the lookup cannot be answered, keep the
chunks for vacuum to reclaim; only a confirmed absence still cleans up.

* s3api: regression tests for ambiguous CreateEntry outcomes

Covers the three post-create-failure cases in putToFiler: the entry
landed despite the error (treat as success, keep chunks), the entry is
confirmed absent (delete orphans), and the outcome is unverifiable
(keep chunks, return error).

* volume: count reads served from deleted needles

A readDeleted read succeeding on a tombstoned needle is the signal that
metadata still points at deleted data. Count it under a
readDeletedNeedle handler label in both the Go and Rust volume servers
so the condition is visible before vacuum turns it into a 404.

* s3api: never delete chunks on an ambiguous create error

Review feedback on the first fix showed verification could still go
wrong in both directions: a stale or lagged lookup could report
not-found for a committed entry, a prefix object stores its chunks on a
directory entry, and filer-side manifestization rewrites the top-level
chunk ids the comparison relied on.

Rework the rule so the outcome classes are asymmetric:

- A transport-level error (anything filerErrorToS3Error maps to a
  retryable 503) is ambiguous and never deletes chunks; the lookup can
  only upgrade the write to success.
- Any other error is a definitive filer refusal and still cleans up.

confirmCreateLanded asks the write owner first, resolves the stored
entry through chunk manifests, requires an exact match of the uploaded
file ids, and on success runs the finalize callback the failed create
skipped (under the object write lock, with the same rmObject undo the
create path uses). Zero-chunk writes stay ambiguous since they cannot
be told apart by chunks.

* s3api: cover definitive refusals and stale entries in put tests

The confirmed-failure case now uses a definitive refusal so it still
exercises orphan cleanup, and a new case keeps chunks when the stored
entry belongs to an older object rather than this PUT.

* volume: count deleted-needle reads once per request

Streamed Go reads ran the deleted check in readNeedle and again in
readNeedleDataInto, and non-streamed Rust reads in stream_info and the
full-read fallback, double-counting one request. Count at the single
entry probe each implementation takes per GET: readNeedle in Go,
read_needle_stream_info in Rust.

* s3api: run recovered-write rollback under the object lock

Two follow-ups from review: ResolveChunkManifest returns traversed
manifest blobs in its manifestChunks output, so requiring it empty
rejected every manifestized landing; and the rmObject undo ran after
the object write lock was released, so a concurrent newer write could
be deleted between finalize failure and rollback. Compare only the
resolved data chunks and keep the undo inside the lock.

* s3api: verify, finalize and roll back recovered creates in one lock

A lookup done before the object write lock let a concurrent PUT replace
the entry between the chunk comparison and the finalize/rollback
section, so a failed afterCreate could rmObject a newer write. Run the
owner lookup, manifest resolution, chunk comparison, afterCreate and
the conditional undo inside a single withObjectWriteLock section.
2026-09-17 19:58:49 -07:00

309 lines
9.4 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 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()
// 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 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)))
defer mem.Free(buf)
// read needle data
crc := needle.CRC(0)
for x := offset; x < offset+size; x += int64(len(buf)) {
if readOption.HasSlowRead {
v.dataFileAccessLock.RLock()
}
// 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])
// Note: CRC validation happens after the loop completes (see below)
// to avoid performance overhead in the hot read path
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 offset == 0 && size == int64(n.DataSize) && (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.
stats.VolumeServerHandlerCounter.WithLabelValues(stats.ErrorCRC).Inc()
return fmt.Errorf("ReadNeedleData checksum %v expected %v for Needle: %v,%v", crc, n.Checksum, v.Id, n)
}
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) {
v.dataFileAccessLock.RLock()
defer v.dataFileAccessLock.RUnlock()
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)
}
offset += NeedleHeaderSize + rest
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
}