Files
seaweedfs/weed/storage/volume_write.go
T
Chris Lu 2ffa696809 fix(volume): handle faulty storage media (Go + Rust) (#11233)
* fix(volume): track EC shard read errors and unmount on faulty media

Extract the volume EIO tracker into a reusable IoErrorTracker and add the
same tracking to EcVolume. Sustained EIO on .ecx lookups or .ecd shard
reads now unmounts the EC volume in the heartbeat (without deleting
files) so the master re-replicates from healthy peers, mirroring the
existing volume replica quarantine.

Closes #11227 (EC shard unmount).

* rust(volume): mirror EC shard read error tracking and unmount

Add EIO tracking to the Rust EcVolume mirroring Go: a streak counter
with IO_ERROR_TOLERANCE, a sticky quarantine flag, and unmount (not
file deletion) in the heartbeat so the master re-replicates from
healthy peers.

* feat(metrics): expose storage IO error counter and quarantine gauge

Add a storage_io_error_total counter incremented on every EIO recorded
by the volume or EC shard tracker, and an io_quarantine gauge labelled
by kind (volume/ec_shard) reflecting the count of replicas suppressed
in the heartbeat. Mirrored in Go and Rust.

* feat(healthz): report 503 when local replicas are IO-quarantined

Add Store.HasIoQuarantine (Go) / Store::has_io_quarantine (Rust) and
have /healthz return 503 when any local volume or EC shard is
quarantined due to sustained storage-media EIO, so a load balancer
can drain a server whose underlying media is faulty. Mirrored in Go
and Rust.

* fix(volume): keep quarantined EC volumes in memory and reset EIO on success

Address review feedback: instead of unloading quarantined EC volumes
(which discards the quarantine state healthz needs), keep them in
memory and just skip them from heartbeat reporting, mirroring the
regular volume quarantine. Also clear the EIO streak on successful
.ecx reads in Rust so a transient error does not accumulate, and add
an ec_shard label to the io_quarantine gauge in both Go and Rust.

* fix(volume): exclude quarantined EC shards from heartbeat and add Rust volume tolerance

Address review feedback:
- Filter quarantined EC volumes from CollectErasureCodingHeartbeat
  (Go) and collect_ec_shard_delta_messages / collect_live_ec_shards
  (Rust) so the master stops advertising faulty shards and
  re-replicates from healthy peers.
- Add consecutive EIO count and sticky quarantine to the Rust
  regular Volume, mirroring Go IoErrorTracker: a single EIO no
  longer deletes the replica; the heartbeat quarantines after the
  tolerance threshold and keeps the volume in memory.
- Use the quarantine flag (not last_io_error) in has_io_quarantine
  so /healthz reflects sustained, not transient, failures.

* fix(volume): make Rust quarantined volumes read-only and wire recovery

Address Devin review:
- Set no_write_or_delete on Rust volumes when quarantined in the
  heartbeat, so cached or direct clients cannot mutate a faulty
  replica after the master removes it (mirrors Go).
- Wire reset_io_error_state into Volume::set_writable so an operator
  making a volume writable again clears the sticky quarantine and
  the volume re-enters heartbeat rotation.

* fix(volume): clear EC quarantine on shard re-mount for operator recovery

Address Greptile review: re-mounting EC shards (Go loadEcShardWithIdxDir
/ Rust mount_ec_shards_with_idx_dir) now calls ResetIoErrorState on the
existing EcVolume, giving operators a documented recovery path that
clears the sticky quarantine and returns the EC volume to heartbeat
rotation. Mirrored in Go and Rust.

* fix(volume): do not clear EC quarantine on routine shard mounts

Address review feedback: clearing the EC IO quarantine on every mount
(including duplicate, retry, sibling-shard, and reconciliation mounts)
is too aggressive and can re-advertise known-bad shards before the
storage media has been validated. Remove the automatic reset from the
mount path; quarantine clears naturally on restart or full unmount
when a fresh EcVolume is created with clean state.

* test(volume): update Rust IO error test for quarantine semantics

The heartbeat now quarantines a volume with sustained EIO (keeps it
mounted, makes it read-only, omits it from heartbeat) instead of
deleting it. Update test_collect_heartbeat_deletes_io_error_volume to
assert the volume stays in the store with no_write_or_delete set, and
update set_last_io_error_for_test to set the consecutive error count
at the tolerance threshold so the test reflects a sustained error.

* fix(volume): reset EIO streak after full write and match Windows media errors

Move the success-side EIO reset from append_needle (after write_all only)
to the end of do_write_request, after flush_dat/flush_idx complete, so a
successful write_all followed by a failed fsync no longer resets the
counter before the EIO is recorded. Repeated fsync EIOs now accumulate
toward the quarantine threshold as intended.

Recognize Windows storage-media failure codes ERROR_CRC (23) and
ERROR_IO_DEVICE (1117) in addition to Unix EIO (errno 5), so quarantined
heartbeat behavior is preserved on Windows. Mirrors the change in both
Go and Rust volume servers.

* fix(volume): preserve checkpoint EIO and clear streak on successful delete

maybe_checkpoint_index now returns whether the checkpoint succeeded;
the success-side EIO reset in do_write_request and do_delete_request
only fires when it did, so a checkpoint media failure is no longer
erased by the unconditional reset that followed it. do_delete_request
also gains the success reset that was lost when append_needle stopped
clearing the streak, so a successful delete still clears an earlier
failure streak.

is_storage_io_error now uses libc::EIO on Unix instead of a hard-coded
5, and the ECX binary-search read path gains a Windows fallback
(seek + read_exact) so the buffer is no longer zeroed on non-Unix
targets.
2026-09-08 21:42:56 -07:00

506 lines
16 KiB
Go

package storage
import (
"bytes"
"errors"
"fmt"
"os"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/storage/backend"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
. "github.com/seaweedfs/seaweedfs/weed/storage/types"
)
var ErrorNotFound = errors.New("not found")
var ErrorDeleted = errors.New("already deleted")
var ErrorSizeMismatch = errors.New("size mismatch")
// isFileUnchanged checks whether this needle to write is same as last one.
// It requires serialized access in the same volume.
func (v *Volume) isFileUnchanged(n *needle.Needle) bool {
if v.Ttl.String() != "" {
return false
}
nv, ok := v.nm.Get(n.Id)
if ok && !nv.Offset.IsZero() && nv.Size.IsValid() {
oldNeedle := new(needle.Needle)
err := oldNeedle.ReadData(v.DataBackend, nv.Offset.ToActualOffset(), nv.Size, v.Version())
if err != nil {
glog.V(0).Infof("Failed to check updated file at offset %d size %d: %v", nv.Offset.ToActualOffset(), nv.Size, err)
return false
}
if oldNeedle.Cookie == n.Cookie && oldNeedle.Checksum == n.Checksum && bytes.Equal(oldNeedle.Data, n.Data) {
n.DataSize = oldNeedle.DataSize
return true
}
}
return false
}
var ErrVolumeNotEmpty = fmt.Errorf("volume not empty")
// Destroy removes everything related to this volume. When keepRemoteData is
// true the cloud-tier object backing the volume is left intact — used by
// moves where another server is taking over the same .vif.
func (v *Volume) Destroy(onlyEmpty bool, keepRemoteData bool) (err error) {
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
if onlyEmpty {
isEmpty, e := v.doIsEmpty()
if e != nil {
err = fmt.Errorf("failed to read isEmpty %v", e)
return
}
if !isEmpty {
err = ErrVolumeNotEmpty
return
}
}
if !v.isCompactionInProgress.CompareAndSwap(false, true) {
err = fmt.Errorf("volume %d is compacting", v.Id)
return
}
v.stopWorker()
if !keepRemoteData {
storageName, storageKey := v.RemoteStorageNameKey()
if v.HasRemoteFile() && storageName != "" && storageKey != "" {
if backendStorage, found := backend.BackendStorages[storageName]; found {
backendStorage.DeleteFile(storageKey)
}
}
}
// A regular volume and an EC volume for the same id share <base>.vif. When
// EC artefacts coexist on this disk (e.g. shards distributed onto a source
// replica before it is deleted), keep the .vif so removing the regular
// volume does not strip the EC volume's info file.
keepVif := v.sharesVifWithEcVolume()
v.doClose()
removeVolumeFiles(v.DataFileName(), keepVif)
removeVolumeFiles(v.IndexFileName(), keepVif)
return
}
// sharesVifWithEcVolume reports whether an EC volume for this volume id lives
// on the same disk, in which case its .vif is the same file as the regular
// volume's and must outlive the regular volume's deletion.
func (v *Volume) sharesVifWithEcVolume() bool {
if v.location == nil {
return false
}
if _, found := v.location.FindEcVolume(v.Id); found {
return true
}
return v.location.HasEcxFileOnDisk(v.Collection, v.Id)
}
func removeVolumeFiles(filename string, keepVif bool) {
// .dat/.idx removals log at V(0) so destructive calls are traceable.
deleteAndLog := func(ext string) {
fullFilename := filename + "." + ext
st, statErr := os.Stat(fullFilename)
err := os.RemoveAll(fullFilename)
if err != nil {
glog.V(0).Infof("failed to remove volume file %s: %s", fullFilename, err)
return
}
if statErr == nil && (ext == "dat" || ext == "idx") {
glog.Infof("removed volume file %s (size=%d)", fullFilename, st.Size())
}
}
deleteAndLog("dat")
deleteAndLog("idx")
if !keepVif {
deleteAndLog("vif")
}
// sorted index file
deleteAndLog("sdx")
// compaction
deleteAndLog("cpd")
deleteAndLog("cpx")
// compaction commit marker
deleteAndLog("cpc")
// level db index file
deleteAndLog("ldb")
// redb index file (Rust volume server)
deleteAndLog("rdb")
// marker for damaged or incomplete volume
deleteAndLog("note")
}
// asyncRequestAppend queues a request for the batch worker, starting it on the
// first one. It reports false for a destroyed volume, so the caller writes
// inline rather than wait on a worker that will never answer.
func (v *Volume) asyncRequestAppend(request *needle.AsyncRequest) bool {
requests := v.startWorker()
if requests == nil {
return false
}
requests <- request
return true
}
func (v *Volume) syncWrite(n *needle.Needle, checkCookie bool, fsync bool) (offset uint64, size Size, isUnchanged bool, err error) {
// glog.V(4).Infof("writing needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
// A caller can still hold the volume after it was closed or destroyed, which
// leaves both of these nil. Refuse the write rather than dereference them.
if v.nm == nil || v.DataBackend == nil {
return 0, 0, false, fmt.Errorf("volume %d is closed", v.Id)
}
if !fsync {
return v.doWriteRequest(n, checkCookie)
}
end, _, statErr := v.DataBackend.GetStat()
if statErr != nil {
return 0, 0, false, fmt.Errorf("cannot read current volume position: %v", statErr)
}
priorOffset, priorSize, hasPrior := Offset{}, Size(0), false
if nv, found := v.nm.Get(n.Id); found {
priorOffset, priorSize, hasPrior = nv.Offset, nv.Size, true
}
offset, size, isUnchanged, err = v.doWriteRequest(n, checkCookie)
if err != nil {
return
}
if syncErr := v.DataBackend.Sync(); syncErr != nil {
v.checkReadWriteError(syncErr)
if !isUnchanged {
v.rollbackUnflushedWrite(n, offset, end, priorOffset, priorSize, hasPrior)
}
return 0, 0, false, syncErr
}
return
}
// rollbackUnflushedWrite undoes an append whose fsync failed: the bytes are not
// data we can vouch for, so they come back off the .dat and the needle map goes
// back to what it pointed at before, rather than at an offset past the new end.
func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int64, priorOffset Offset, priorSize Size, hasPrior bool) {
if te := v.DataBackend.Truncate(end); te != nil {
glog.V(0).Infof("Failed to truncate %s back to %d with error: %v", v.DataBackend.Name(), end, te)
}
current, found := v.nm.Get(n.Id)
if !found || current.Offset.ToActualOffset() != int64(offset) {
// doWriteRequest kept an existing mapping at a higher offset
return
}
var err error
if hasPrior {
err = v.nm.Put(n.Id, priorOffset, priorSize)
} else {
err = v.nm.Delete(n.Id, ToOffset(int64(offset)))
}
if err != nil {
glog.V(0).Infof("Failed to roll back the index of needle %d in volume %d: %v", n.Id, v.Id, err)
}
}
// writeNeedle2 appends a needle. A durable write normally goes through the
// async batch worker, which fsyncs once for the whole batch; while the server
// is stopping the worker is winding down, so it is flushed inline instead. Both
// paths only return once the .dat is on disk.
func (v *Volume) writeNeedle2(n *needle.Needle, checkCookie bool, fsync bool, isStopping bool) (offset uint64, size Size, isUnchanged bool, err error) {
// glog.V(4).Infof("writing needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
if n.Ttl == needle.EMPTY_TTL && v.Ttl != needle.EMPTY_TTL {
n.SetHasTtl()
n.Ttl = v.Ttl
}
if !fsync || isStopping {
return v.syncWrite(n, checkCookie, fsync)
} else {
asyncRequest := needle.NewAsyncRequest(n, true)
// using len(n.Data) here instead of n.Size before n.Size is populated in n.Append()
asyncRequest.ActualSize = needle.GetActualSize(Size(len(n.Data)), v.Version())
if !v.asyncRequestAppend(asyncRequest) {
return v.syncWrite(n, checkCookie, fsync)
}
offset, _, isUnchanged, err = asyncRequest.WaitComplete()
return
}
}
func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool) (offset uint64, size Size, isUnchanged bool, err error) {
// glog.V(4).Infof("writing needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
if v.isFileUnchanged(n) {
size = Size(n.DataSize)
isUnchanged = true
return
}
// check whether existing needle cookie matches
nv, ok := v.nm.Get(n.Id)
if ok {
existingNeedle, _, _, existingNeedleReadErr := needle.ReadNeedleHeader(v.DataBackend, v.Version(), nv.Offset.ToActualOffset())
if existingNeedleReadErr != nil {
err = fmt.Errorf("reading existing needle: %w", existingNeedleReadErr)
return
}
if n.Cookie == 0 && !checkCookie {
// this is from batch deletion, and read back again when tailing a remote volume
// which only happens when checkCookie == false and fsync == false
n.Cookie = existingNeedle.Cookie
}
if existingNeedle.Cookie != n.Cookie {
glog.V(0).Infof("write cookie mismatch: existing %s, new %s",
needle.NewFileIdFromNeedle(v.Id, existingNeedle), needle.NewFileIdFromNeedle(v.Id, n))
err = fmt.Errorf("mismatching cookie %x", n.Cookie)
return
}
}
// append to dat file
n.UpdateAppendAtNs(v.lastAppendAtNs)
var actualSize int64
offset, size, actualSize, err = n.Append(v.DataBackend, v.Version())
v.checkReadWriteError(err)
if err != nil {
err = fmt.Errorf("append to volume %d size %d actualSize %d: %v", v.Id, size, actualSize, err)
return
}
v.lastAppendAtNs = n.AppendAtNs
// add to needle map
if !ok || uint64(nv.Offset.ToActualOffset()) < offset {
if err = v.nm.Put(n.Id, ToOffset(int64(offset)), n.Size); err != nil {
err = fmt.Errorf("index needle %d of volume %d at offset %d: %w", n.Id, v.Id, offset, err)
glog.V(0).Info(err)
}
}
if v.lastModifiedTsSeconds < n.LastModified {
v.lastModifiedTsSeconds = n.LastModified
}
return
}
func (v *Volume) syncDelete(n *needle.Needle) (Size, error) {
// glog.V(4).Infof("delete needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
if v.nm == nil {
return 0, nil
}
return v.doDeleteRequest(n)
}
func (v *Volume) deleteNeedle2(n *needle.Needle) (Size, error) {
// todo: delete info is always appended no fsync, it may need fsync in future
fsync := false
if !fsync {
return v.syncDelete(n)
} else {
asyncRequest := needle.NewAsyncRequest(n, false)
asyncRequest.ActualSize = needle.GetActualSize(0, v.Version())
if !v.asyncRequestAppend(asyncRequest) {
return v.syncDelete(n)
}
_, size, _, err := asyncRequest.WaitComplete()
return Size(size), err
}
}
func (v *Volume) doDeleteRequest(n *needle.Needle) (Size, error) {
glog.V(4).Infof("delete needle %s", needle.NewFileIdFromNeedle(v.Id, n).String())
nv, ok := v.nm.Get(n.Id)
// fmt.Println("key", n.Id, "volume offset", nv.Offset, "data_size", n.Size, "cached size", nv.Size)
if ok && !nv.Size.IsDeleted() {
var offset uint64
var err error
size := nv.Size
if !v.HasRemoteFile() {
n.Data = nil
n.UpdateAppendAtNs(v.lastAppendAtNs)
offset, _, _, err = n.Append(v.DataBackend, v.Version())
v.checkReadWriteError(err)
if err != nil {
return size, err
}
}
v.lastAppendAtNs = n.AppendAtNs
if err = v.nm.Delete(n.Id, ToOffset(int64(offset))); err != nil {
return size, err
}
return size, err
}
return 0, nil
}
// startWorker returns the volume's batch-write channel, creating it and its
// goroutine on first use, and nil once stopWorker has run.
func (v *Volume) startWorker() chan *needle.AsyncRequest {
v.asyncWorkerLock.Lock()
defer v.asyncWorkerLock.Unlock()
if v.asyncWorkerClosed {
return nil
}
if v.asyncRequestsChan != nil {
return v.asyncRequestsChan
}
requests := make(chan *needle.AsyncRequest, 128)
v.asyncRequestsChan = requests
go func() {
chanClosed := false
for {
// chan closed. go thread will exit
if chanClosed {
break
}
currentRequests := make([]*needle.AsyncRequest, 0, 128)
currentBytesToWrite := int64(0)
for {
request, ok := <-requests
// volume may be closed
if !ok {
chanClosed = true
break
}
if MaxPossibleVolumeSize < v.ContentSize()+uint64(currentBytesToWrite+request.ActualSize) {
request.Complete(0, 0, false,
fmt.Errorf("volume size limit %d exceeded! current size is %d", MaxPossibleVolumeSize, v.ContentSize()))
break
}
currentRequests = append(currentRequests, request)
currentBytesToWrite += request.ActualSize
// submit at most 4M bytes or 128 requests at one time to decrease request delay.
// it also need to break if there is no data in channel to avoid io hang.
if currentBytesToWrite >= 4*1024*1024 || len(currentRequests) >= 128 || len(requests) == 0 {
break
}
}
if len(currentRequests) == 0 {
continue
}
v.dataFileAccessLock.Lock()
end, e := int64(0), error(nil)
if v.nm == nil || v.DataBackend == nil {
e = fmt.Errorf("volume %d is closed", v.Id)
} else {
end, _, e = v.DataBackend.GetStat()
}
if e != nil {
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Complete(0, 0, false,
fmt.Errorf("cannot read current volume position: %v", e))
}
v.dataFileAccessLock.Unlock()
continue
}
for i := 0; i < len(currentRequests); i++ {
if currentRequests[i].IsWriteRequest {
offset, size, isUnchanged, err := v.doWriteRequest(currentRequests[i].N, true)
currentRequests[i].UpdateResult(offset, uint64(size), isUnchanged, err)
} else {
size, err := v.doDeleteRequest(currentRequests[i].N)
currentRequests[i].UpdateResult(0, uint64(size), false, err)
}
}
// if sync error, data is not reliable, we should mark the completed request as fail and rollback
if err := v.DataBackend.Sync(); err != nil {
// todo: this may generate dirty data or cause data inconsistent, may be weed need to panic?
if te := v.DataBackend.Truncate(end); te != nil {
glog.V(0).Infof("Failed to truncate %s back to %d with error: %v", v.DataBackend.Name(), end, te)
}
for i := 0; i < len(currentRequests); i++ {
if currentRequests[i].IsSucceed() {
currentRequests[i].UpdateResult(0, 0, false, err)
}
}
}
for i := 0; i < len(currentRequests); i++ {
currentRequests[i].Submit()
}
v.dataFileAccessLock.Unlock()
}
}()
return requests
}
// stopWorker closes the batch-write channel so the worker drains what is queued
// and exits. It stays closed: a destroyed volume takes no more writes.
func (v *Volume) stopWorker() {
v.asyncWorkerLock.Lock()
defer v.asyncWorkerLock.Unlock()
if v.asyncWorkerClosed {
return
}
v.asyncWorkerClosed = true
if v.asyncRequestsChan != nil {
close(v.asyncRequestsChan)
v.asyncRequestsChan = nil
}
}
func (v *Volume) WriteNeedleBlob(needleId NeedleId, needleBlob []byte, size Size) error {
v.dataFileAccessLock.Lock()
defer v.dataFileAccessLock.Unlock()
// nm.Put on a read-only volume fails only after the blob is appended to .dat.
if v.IsReadOnly() {
return fmt.Errorf("volume %d is read only", v.Id)
}
// size indexes the needle and places the v3 append timestamp, so a caller using
// the payload-only DataSize corrupts both, silently until the needle is read back.
if len(needleBlob) < NeedleHeaderSize {
return fmt.Errorf("needle %d blob of %d bytes is shorter than a needle header", needleId, len(needleBlob))
}
var blobHeader needle.Needle
blobHeader.ParseNeedleHeader(needleBlob)
if blobHeader.Size != size {
return fmt.Errorf("needle %d size %d does not match its blob header size %d", needleId, size, blobHeader.Size)
}
if MaxPossibleVolumeSize < v.nm.ContentSize()+uint64(len(needleBlob)) {
return fmt.Errorf("volume size limit %d exceeded! current size is %d", MaxPossibleVolumeSize, v.nm.ContentSize())
}
nv, ok := v.nm.Get(needleId)
if ok && nv.Size == size {
oldNeedle := new(needle.Needle)
err := oldNeedle.ReadData(v.DataBackend, nv.Offset.ToActualOffset(), nv.Size, v.Version())
if err == nil {
newNeedle := new(needle.Needle)
err = newNeedle.ReadBytes(needleBlob, nv.Offset.ToActualOffset(), size, v.Version())
if err == nil && oldNeedle.Cookie == newNeedle.Cookie && oldNeedle.Checksum == newNeedle.Checksum && bytes.Equal(oldNeedle.Data, newNeedle.Data) {
glog.V(0).Infof("needle %v already exists", needleId)
return nil
}
}
}
appendAtNs := needle.GetAppendAtNs(v.lastAppendAtNs)
offset, err := needle.WriteNeedleBlob(v.DataBackend, needleBlob, size, appendAtNs, v.Version())
v.checkReadWriteError(err)
if err != nil {
return err
}
v.lastAppendAtNs = appendAtNs
// add to needle map
if err = v.nm.Put(needleId, ToOffset(int64(offset)), size); err != nil {
err = fmt.Errorf("index needle %d of volume %d at offset %d: %w", needleId, v.Id, offset, err)
glog.V(0).Info(err)
}
return err
}