Files
seaweedfs/weed/server/volume_grpc_erasure_coding_test.go
T
e3e02d3364 [CheckDisk]: implement disk health detection (#9560)
* [CheckDisk][GRPC]: implement MVP for disk health detection, added timeout for new grpc connections

* fix(volume): build disk health check on every platform

setDiskStatus only existed behind the statfs build tag, so disk.go failed
to compile on windows, openbsd, solaris, netbsd and plan9. Move the timeout
wrapper and failure tracking into the shared disk.go and have each platform's
fillInDiskStatus return an error, so every platform gets the same protection
from a stuck filesystem.

Also restore the uint64(fs.Bavail) cast: Bavail is int64 on freebsd, so the
unguarded multiply broke the freebsd build.

* fix(volume): keep one outstanding statfs probe per disk

A stuck statfs used to leave isChecking cleared by the timeout path, so the
next check spawned another goroutine while the previous one was still blocked
in the syscall, leaking one goroutine per minute on a hung disk. Clear the
flag only when statfs returns and treat an overlapping check as a failure, so
a hung filesystem keeps a single outstanding probe and still gets reported.

* fix(volume): assume disk available until the first health check

isDiskAvailable defaulted to false, and CollectHeartbeat skips locations that
are not available. A freshly started volume server would therefore omit every
volume from its first heartbeats until the async CheckDiskSpace ran, so the
master could briefly treat all of them as missing.

* fix(volume): label the disk error metric by data directory

The new gauge tagged the series with IdxDirectory while every neighbouring
resource gauge uses Directory, so the error series would not line up with them
in dashboards. Also log the underlying error instead of a generic message.

* test(volume): cover disk health success and repeated-failure paths

* fix(volume): make a healthy disk the zero-value default

Track the disk as isDiskUnavailable instead of isDiskAvailable so the safe
state is the zero value, matching isDiskSpaceLow. CollectHeartbeat only skips a
location once a check has actively marked it unavailable, so any DiskLocation
built without running CheckDiskSpace (tests, future call sites) still reports
its volumes instead of silently dropping them.

* feat(disk): detect degraded disks using IO latency probes

* feat(stats): introduce configurable disk I/O health probe with EWMA-based latency detection

* feat(disk): replace EWMA with sliding window algorithm for disk health detection and added user-friendly options

* feat(disk): improve disk health probing and recovery

* feat(volume): configure disk health checks via volume.toml

* fix(volume): Remove disk IO probe CLI options

---------

Co-authored-by: ptukha <ptukha@tochka.com>
Co-authored-by: Chris Lu <chris.lu@gmail.com>
2026-06-02 09:02:05 -07:00

206 lines
6.3 KiB
Go

package weed_server
import (
"context"
"os"
"path/filepath"
"sort"
"testing"
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
"github.com/seaweedfs/seaweedfs/weed/stats"
"github.com/seaweedfs/seaweedfs/weed/storage"
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
"github.com/seaweedfs/seaweedfs/weed/storage/types"
"github.com/seaweedfs/seaweedfs/weed/storage/volume_info"
"github.com/seaweedfs/seaweedfs/weed/util"
)
func TestCheckEcVolumeStatusCountOnlyDataShards(t *testing.T) {
tempDir := t.TempDir()
dataDir := filepath.Join(tempDir, "data")
idxDir := filepath.Join(tempDir, "idx")
if err := os.MkdirAll(dataDir, 0o755); err != nil {
t.Fatalf("mkdir data dir: %v", err)
}
if err := os.MkdirAll(idxDir, 0o755); err != nil {
t.Fatalf("mkdir idx dir: %v", err)
}
baseName := "7"
filesToCreate := []string{
filepath.Join(dataDir, baseName+".ec00"),
filepath.Join(dataDir, baseName+".ec09"),
filepath.Join(dataDir, baseName+".ec13"),
filepath.Join(idxDir, baseName+".ecx"),
filepath.Join(idxDir, baseName+".ecj"),
filepath.Join(idxDir, baseName+".idx"),
}
for _, fileName := range filesToCreate {
if err := os.WriteFile(fileName, []byte("x"), 0o644); err != nil {
t.Fatalf("create %s: %v", fileName, err)
}
}
location := &storage.DiskLocation{
Directory: dataDir,
IdxDirectory: idxDir,
}
hasEcxFile, hasIdxFile, shardCount, err := checkEcVolumeStatus(baseName, location)
if err != nil {
t.Fatalf("checkEcVolumeStatus: %v", err)
}
if !hasEcxFile {
t.Fatalf("expected hasEcxFile=true")
}
if !hasIdxFile {
t.Fatalf("expected hasIdxFile=true")
}
if shardCount != 3 {
t.Fatalf("expected shardCount=3, got %d", shardCount)
}
}
// TestVolumeEcShardsInfo_AggregatesAcrossDisks pins the multi-disk path:
// when a volume server mounts EC shards for the same volume on more than
// one disk (each disk holds its own EcVolume entry — Store.FindEcVolume
// returns only the first), VolumeEcShardsInfo used to report shards from
// a single disk. The ec.encode verification step (verifyEcShardsBeforeDelete)
// then refused to delete the source volume because the union across
// servers fell short of dataShards + parityShards. The handler must walk
// every DiskLocation so the response covers every shard the server holds.
func TestVolumeEcShardsInfo_AggregatesAcrossDisks(t *testing.T) {
tempDir := t.TempDir()
dir0 := filepath.Join(tempDir, "disk0")
dir1 := filepath.Join(tempDir, "disk1")
for _, d := range []string{dir0, dir1} {
if err := os.MkdirAll(d, 0o755); err != nil {
t.Fatalf("mkdir %s: %v", d, err)
}
}
const collection = "ec-multi-disk-info"
vid := needle.VolumeId(42)
const dataShards, parityShards = 10, 4
const datSize int64 = 10 * 1024 * 1024
// Two shards on disk0, two on disk1. The .ecx / .ecj / .vif live on
// disk0 so each disk's EcVolume can open the index files via the
// cross-disk fallback in NewEcVolume.
shardsOnDisk0 := []erasure_coding.ShardId{0, 5}
shardsOnDisk1 := []erasure_coding.ShardId{7, 12}
diskIOProbeConfig := stats.DefaultDiskIOProbeConfig()
store := storage.NewStore(nil, "localhost", 8080, 18080, "http://localhost:8080", "store-id",
[]string{dir0, dir1},
[]int32{100, 100},
[]util.MinFreeSpace{{}, {}},
"",
storage.NeedleMapInMemory,
[]types.DiskType{types.HardDriveType, types.HardDriveType},
nil,
3,
diskIOProbeConfig,
)
done := make(chan struct{})
go func() {
for {
select {
case <-store.NewEcShardsChan:
case <-store.NewVolumesChan:
case <-store.DeletedVolumesChan:
case <-store.DeletedEcShardsChan:
case <-store.StateUpdateChan:
case <-done:
return
}
}
}()
t.Cleanup(func() {
store.Close()
close(done)
})
base0 := erasure_coding.EcShardFileName(collection, dir0, int(vid))
base1 := erasure_coding.EcShardFileName(collection, dir1, int(vid))
// .ecx, .ecj, .vif live on disk0. NewEcVolume on disk1 falls back to
// disk0's idx dir. The .ecx needs >0 bytes so HasEcxFileOnDisk does
// not treat it as the corrupt-stub case; a single zero entry is the
// smallest valid index file (WalkIndex iterates one zero-sized needle).
if err := os.WriteFile(base0+".ecx", make([]byte, 16), 0o644); err != nil {
t.Fatalf("write .ecx: %v", err)
}
if err := os.WriteFile(base0+".ecj", nil, 0o644); err != nil {
t.Fatalf("write .ecj: %v", err)
}
if err := volume_info.SaveVolumeInfo(base0+".vif", &volume_server_pb.VolumeInfo{
Version: uint32(needle.Version3),
DatFileSize: datSize,
EcShardConfig: &volume_server_pb.EcShardConfig{
DataShards: dataShards,
ParityShards: parityShards,
},
}); err != nil {
t.Fatalf("save .vif: %v", err)
}
plant := func(base string, shardId erasure_coding.ShardId) {
t.Helper()
f, err := os.Create(base + erasure_coding.ToExt(int(shardId)))
if err != nil {
t.Fatalf("create shard %d: %v", shardId, err)
}
// MountEcShards does not validate shard size, so any non-empty
// truncate avoids the zero-byte ignore branch in loadAllEcShards.
if err := f.Truncate(1); err != nil {
f.Close()
t.Fatalf("truncate shard %d: %v", shardId, err)
}
f.Close()
}
for _, sid := range shardsOnDisk0 {
plant(base0, sid)
}
for _, sid := range shardsOnDisk1 {
plant(base1, sid)
}
for _, sid := range append([]erasure_coding.ShardId{}, append(shardsOnDisk0, shardsOnDisk1...)...) {
if err := store.MountEcShards(collection, vid, sid, ""); err != nil {
t.Fatalf("MountEcShards %d.%d: %v", vid, sid, err)
}
}
vs := &VolumeServer{store: store}
resp, err := vs.VolumeEcShardsInfo(context.Background(), &volume_server_pb.VolumeEcShardsInfoRequest{
VolumeId: uint32(vid),
})
if err != nil {
t.Fatalf("VolumeEcShardsInfo: %v", err)
}
gotShardIds := make([]int, 0, len(resp.GetEcShardInfos()))
for _, info := range resp.GetEcShardInfos() {
if info.GetVolumeId() != uint32(vid) {
t.Errorf("EcShardInfo VolumeId=%d, want %d", info.GetVolumeId(), vid)
}
gotShardIds = append(gotShardIds, int(info.GetShardId()))
}
sort.Ints(gotShardIds)
want := []int{0, 5, 7, 12}
if len(gotShardIds) != len(want) {
t.Fatalf("VolumeEcShardsInfo returned %d shards (ids=%v), want %d (ids=%v)",
len(gotShardIds), gotShardIds, len(want), want)
}
for i, sid := range want {
if gotShardIds[i] != sid {
t.Fatalf("VolumeEcShardsInfo shard ids=%v, want %v", gotShardIds, want)
}
}
}