From c858e01a097dcf3383b9a0d838fae61533e098bc Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Fri, 28 Aug 2026 15:55:12 -0700 Subject: [PATCH] ec: split the shard-interval recovery into a gather and a rebuild (#11005) * ec: split the shard-interval recovery into a gather and a rebuild Recovering an interval is now one function doing the local seeding, the waved peer fetch, the shard accounting and the Reed-Solomon rebuild, under a memory budget. Splitting the gather from the rebuild makes the rebuild a plain function over a set of intervals, which is testable on its own and reusable by the parity checks a full scrub wants. The rebuild refuses a parity target, and the caller checks that before the gather so a doomed target costs no fan-out. ReconstructData rebuilds data shards only, so asking it for a parity shard returned no error and left the slot nil, and the caller copied that out as a successful read of zeroes. Only data shard ids reach here today, so this is a guard, not a live fix. Claude-Session: https://claude.ai/code/session_014yMNebkUjSbx9sfUCWJJtq * ec: rebuild only the EC shard the read asked for ReconstructData rebuilds every missing data shard. The gather stops as soon as DataShards intervals are in hand, so on a distributed volume it routinely finishes holding parity where data is missing -- and each of those data shards is then rebuilt into an interval-sized buffer, decoded, and never read. Ask for the one shard the read needs. The budget covers it now too: DataShards gathered plus the one the rebuild allocates. It never covered the rebuild's output, and with ReconstructData that output was up to ParityShards buffers. The required mask is Total() long rather than DataShards. reedsolomon documents both lengths, but its presence scan walks every shard and indexes the short mask past its end, so the documented short form panics whenever a parity shard is absent - which here it usually is. Claude-Session: https://claude.ai/code/session_014yMNebkUjSbx9sfUCWJJtq --- weed/storage/store_ec.go | 133 +++++++++++++++------- weed/storage/store_ec_reconstruct_test.go | 118 +++++++++++++++++++ 2 files changed, 207 insertions(+), 44 deletions(-) create mode 100644 weed/storage/store_ec_reconstruct_test.go diff --git a/weed/storage/store_ec.go b/weed/storage/store_ec.go index dc1cbc49a..0abd9fc43 100644 --- a/weed/storage/store_ec.go +++ b/weed/storage/store_ec.go @@ -747,38 +747,11 @@ func (s *Store) doReadRemoteEcShardInterval(sourceDataNode pb.ServerAddress, nee return } -func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, shardIdToRecover erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) { - glog.V(3).Infof("recover ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover) - - // Reconstruct with the volume's OWN EC ratio (loaded from its .vif), not the - // build default, so a custom-ratio volume (e.g. 9+3) is decoded with the matrix - // that actually produced its shards -- decoding it as 10+4 would corrupt the - // recovered bytes. In OSS the ratio is always 10+4, so this is a no-op. - ecCtx := ecVolume.ECContext - if ecCtx == nil { - ecCtx = erasure_coding.NewDefaultECContext(ecVolume.Collection, ecVolume.VolumeId) - } - enc, err := reedsolomon.New(ecCtx.DataShards, ecCtx.ParityShards) - if err != nil { - return 0, false, fmt.Errorf("failed to create encoder: %w", err) - } - - // Charge the buffers this recovery is about to hold against the budget, so a - // burst of them queues here rather than on the heap. - weight := int64(len(buf)) * int64(ecCtx.DataShards) - if weight > ecRecoverBudget { - // An interval whose fan-out outgrows the whole budget takes all of it and - // so runs alone, rather than blocking forever on an acquire that can never - // succeed. The cap is then one such recovery, not a burst of them. - weight = ecRecoverBudget - } - if err = ecRecoverSem.Acquire(context.Background(), weight); err != nil { - return 0, false, err - } - defer ecRecoverSem.Release(weight) - - // Use MaxShardCount to support custom EC ratios up to 32 shards - bufs := make([][]byte, erasure_coding.MaxShardCount) +// gatherEcShardIntervals collects the same interval from every shard but shardIdToRecover, +// returning one buffer per shard and nil where the shard could not be gathered. It stops +// once DataShards of them are in hand: reconstruction consumes no more than that. +func (s *Store) gatherEcShardIntervals(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, ecCtx *erasure_coding.ECContext, shardIdToRecover erasure_coding.ShardId, size int, offset int64) (shardIntervals [][]byte, isDeleted bool) { + shardIntervals = make([][]byte, ecCtx.Total()) // A shard this server already holds costs no round trip and no peer buffer, // so seed those before asking peers for the rest. @@ -790,12 +763,12 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum if _, _, found := s.FindEcVolumeWithShard(ecVolume.VolumeId, shardId); !found { continue } - data := make([]byte, len(buf)) + data := make([]byte, size) if localErr := s.readLocalEcShardInterval(ecVolume, shardId, data, offset); localErr != nil { glog.V(3).Infof("recover: read local ec shard %d.%d: %v", ecVolume.VolumeId, shardId, localErr) continue } - bufs[shardId] = data + shardIntervals[shardId] = data available++ } @@ -805,7 +778,7 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum for shardId, locations := range ecVolume.ShardLocations { // skip the shard being recovered, one already seeded locally, or an empty shard - if shardId == shardIdToRecover || int(shardId) >= ecCtx.Total() || bufs[shardId] != nil { + if shardId == shardIdToRecover || int(shardId) >= ecCtx.Total() || shardIntervals[shardId] != nil { continue } if len(locations) == 0 { @@ -834,7 +807,7 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum wg.Add(1) go func() { defer wg.Done() - data := make([]byte, len(buf)) + data := make([]byte, size) nRead, isDeleted, readErr := s.readRemoteEcShardInterval(locations, needleId, ecVolume.VolumeId, shardId, data, offset, ecVolume.EncodeTsNs) if readErr != nil { glog.V(3).Infof("recover: readRemoteEcShardInterval %d.%d %d bytes from %+v: %v", ecVolume.VolumeId, shardId, nRead, locations, readErr) @@ -843,8 +816,8 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum if isDeleted { isDeletedFlag.Store(true) } - if nRead == len(buf) { - bufs[shardId] = data + if nRead == size { + shardIntervals[shardId] = data fetched.Add(1) } }() @@ -857,13 +830,37 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum break } } - is_deleted = isDeletedFlag.Load() + + return shardIntervals, isDeletedFlag.Load() +} + +// checkEcShardRebuildable rejects a parity target. ReconstructData rebuilds data shards +// only, so it leaves a parity slot nil and reports no error, and the caller would copy +// that out as a zero-filled buffer and report a successful read. +func checkEcShardRebuildable(ecVolume *erasure_coding.EcVolume, ecCtx *erasure_coding.ECContext, shardIdToRecover erasure_coding.ShardId) error { + if int(shardIdToRecover) >= ecCtx.DataShards { + return fmt.Errorf("cannot reconstruct shard %d.%d: only data shards can be rebuilt, %d of %d are parity", + ecVolume.VolumeId, shardIdToRecover, ecCtx.ParityShards, ecCtx.Total()) + } + return nil +} + +// reconstructEcShardInterval rebuilds one shard's interval in place from the others. +func reconstructEcShardInterval(ecVolume *erasure_coding.EcVolume, ecCtx *erasure_coding.ECContext, shardIntervals [][]byte, shardIdToRecover erasure_coding.ShardId) error { + if err := checkEcShardRebuildable(ecVolume, ecCtx, shardIdToRecover); err != nil { + return err + } + + enc, err := reedsolomon.New(ecCtx.DataShards, ecCtx.ParityShards) + if err != nil { + return fmt.Errorf("failed to create encoder: %w", err) + } // Count and log available shards for diagnostics availableShards := make([]erasure_coding.ShardId, 0, ecCtx.Total()) missingShards := make([]erasure_coding.ShardId, 0, ecCtx.ParityShards+1) for shardId := 0; shardId < ecCtx.Total(); shardId++ { - if bufs[shardId] != nil { + if shardIntervals[shardId] != nil { availableShards = append(availableShards, erasure_coding.ShardId(shardId)) } else { missingShards = append(missingShards, erasure_coding.ShardId(shardId)) @@ -876,19 +873,67 @@ func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolum len(missingShards), missingShards) if len(availableShards) < ecCtx.DataShards { - return 0, is_deleted, fmt.Errorf("cannot recover shard %d.%d: only %d shards available %v, need at least %d (missing: %v)", + return fmt.Errorf("cannot recover shard %d.%d: only %d shards available %v, need at least %d (missing: %v)", ecVolume.VolumeId, shardIdToRecover, len(availableShards), availableShards, ecCtx.DataShards, missingShards) } - if err = enc.ReconstructData(bufs[:ecCtx.Total()]); err != nil { - return 0, is_deleted, fmt.Errorf("failed to reconstruct data for shard %d.%d with %d available shards %v: %w", + // Rebuild only what was asked for. ReconstructData rebuilds every missing data + // shard, and a gather that stopped at DataShards can leave up to ParityShards of + // them missing -- an interval-sized buffer and a decode each, discarded unread. + // The mask is Total() long, not DataShards: reedsolomon documents both lengths + // but indexes the short one past its end when a parity shard is absent, which + // here it usually is. + required := make([]bool, ecCtx.Total()) + required[shardIdToRecover] = true + if err := enc.ReconstructSome(shardIntervals, required); err != nil { + return fmt.Errorf("failed to reconstruct data for shard %d.%d with %d available shards %v: %w", ecVolume.VolumeId, shardIdToRecover, len(availableShards), availableShards, err) } + + return nil +} + +func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, shardIdToRecover erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) { + glog.V(3).Infof("recover ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover) + + // Reconstruct with the volume's OWN EC ratio (loaded from its .vif), not the + // build default, so a custom-ratio volume (e.g. 9+3) is decoded with the matrix + // that actually produced its shards -- decoding it as 10+4 would corrupt the + // recovered bytes. In OSS the ratio is always 10+4, so this is a no-op. + ecCtx := ecVolume.ECContext + if ecCtx == nil { + ecCtx = erasure_coding.NewDefaultECContext(ecVolume.Collection, ecVolume.VolumeId) + } + // checked before the gather: a doomed target should not cost a fan-out, nor drop a + // sibling's location through forgetShardId on the way to failing + if err := checkEcShardRebuildable(ecVolume, ecCtx, shardIdToRecover); err != nil { + return 0, false, err + } + + // Charge the buffers this recovery is about to hold against the budget, so a + // burst of them queues here rather than on the heap: DataShards gathered, plus + // the one the rebuild allocates for the shard it recreates. + weight := int64(len(buf)) * int64(ecCtx.DataShards+1) + if weight > ecRecoverBudget { + // An interval whose fan-out outgrows the whole budget takes all of it and + // so runs alone, rather than blocking forever on an acquire that can never + // succeed. The cap is then one such recovery, not a burst of them. + weight = ecRecoverBudget + } + if err = ecRecoverSem.Acquire(context.Background(), weight); err != nil { + return 0, false, err + } + defer ecRecoverSem.Release(weight) + + shardIntervals, is_deleted := s.gatherEcShardIntervals(needleId, ecVolume, ecCtx, shardIdToRecover, len(buf), offset) + if err := reconstructEcShardInterval(ecVolume, ecCtx, shardIntervals, shardIdToRecover); err != nil { + return 0, is_deleted, err + } glog.V(4).Infof("recovered ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover) - copy(buf, bufs[shardIdToRecover]) + copy(buf, shardIntervals[shardIdToRecover]) return len(buf), is_deleted, nil } diff --git a/weed/storage/store_ec_reconstruct_test.go b/weed/storage/store_ec_reconstruct_test.go new file mode 100644 index 000000000..f0ce6c01b --- /dev/null +++ b/weed/storage/store_ec_reconstruct_test.go @@ -0,0 +1,118 @@ +package storage + +import ( + "bytes" + "math/rand" + "strings" + "testing" + + "github.com/klauspost/reedsolomon" + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/storage/needle" +) + +// encodedInterval returns one interval's worth of every shard, Reed-Solomon encoded +// from random data, along with the context describing the ratio. +func encodedInterval(t *testing.T, intervalSize int) ([][]byte, *erasure_coding.ECContext) { + t.Helper() + + ecCtx := erasure_coding.NewDefaultECContext("", 0) + shardIntervals := make([][]byte, ecCtx.Total()) + r := rand.New(rand.NewSource(1)) + for i := range shardIntervals { + shardIntervals[i] = make([]byte, intervalSize) + } + for i := 0; i < ecCtx.DataShards; i++ { + r.Read(shardIntervals[i]) + } + + enc, err := reedsolomon.New(ecCtx.DataShards, ecCtx.ParityShards) + if err != nil { + t.Fatalf("new encoder: %v", err) + } + if err := enc.Encode(shardIntervals); err != nil { + t.Fatalf("encode: %v", err) + } + return shardIntervals, ecCtx +} + +func TestReconstructEcShardIntervalRebuildsDataShard(t *testing.T) { + shardIntervals, ecCtx := encodedInterval(t, 1024) + ecVolume := &erasure_coding.EcVolume{VolumeId: needle.VolumeId(1), ECContext: ecCtx} + + const lost = erasure_coding.ShardId(3) + want := bytes.Clone(shardIntervals[lost]) + shardIntervals[lost] = nil + // drop as many others as parity allows, so the rebuild really goes through parity + for i := ecCtx.Total() - ecCtx.ParityShards + 1; i < ecCtx.Total(); i++ { + shardIntervals[i] = nil + } + + if err := reconstructEcShardInterval(ecVolume, ecCtx, shardIntervals, lost); err != nil { + t.Fatalf("reconstruct: %v", err) + } + if !bytes.Equal(shardIntervals[lost], want) { + t.Fatalf("rebuilt shard %d does not match the encoded bytes", lost) + } +} + +// ReconstructData rebuilds data shards only, leaving a parity slot nil. Reporting that +// as a success would hand the caller a zero-filled buffer as if it had been read. +func TestReconstructEcShardIntervalRejectsParityShard(t *testing.T) { + shardIntervals, ecCtx := encodedInterval(t, 1024) + ecVolume := &erasure_coding.EcVolume{VolumeId: needle.VolumeId(1), ECContext: ecCtx} + + parity := erasure_coding.ShardId(ecCtx.DataShards) + shardIntervals[parity] = nil + + err := reconstructEcShardInterval(ecVolume, ecCtx, shardIntervals, parity) + if err == nil { + t.Fatalf("rebuilding parity shard %d reported success, buffer is %v", parity, shardIntervals[parity]) + } + if !strings.Contains(err.Error(), "only data shards can be rebuilt") { + t.Fatalf("error %q, want it to say parity cannot be rebuilt", err) + } +} + +func TestReconstructEcShardIntervalNeedsDataShardCount(t *testing.T) { + shardIntervals, ecCtx := encodedInterval(t, 1024) + ecVolume := &erasure_coding.EcVolume{VolumeId: needle.VolumeId(1), ECContext: ecCtx} + + // one shard short of what the ratio needs + for i := 0; i <= ecCtx.ParityShards; i++ { + shardIntervals[i] = nil + } + + err := reconstructEcShardInterval(ecVolume, ecCtx, shardIntervals, 0) + if err == nil || !strings.Contains(err.Error(), "need at least") { + t.Fatalf("error %v, want it to report too few shards", err) + } +} + +// A gather that filled up on parity leaves several data shards missing, and only one +// of them is the shard anybody asked for. +func TestReconstructEcShardIntervalRebuildsOnlyTheTarget(t *testing.T) { + shardIntervals, ecCtx := encodedInterval(t, 1024) + ecVolume := &erasure_coding.EcVolume{VolumeId: needle.VolumeId(1), ECContext: ecCtx} + + const lost = erasure_coding.ShardId(0) + want := bytes.Clone(shardIntervals[lost]) + // exactly DataShards left in hand, ParityShards of the data shards missing + spare := []erasure_coding.ShardId{4, 6, 8} + shardIntervals[lost] = nil + for _, sid := range spare { + shardIntervals[sid] = nil + } + + if err := reconstructEcShardInterval(ecVolume, ecCtx, shardIntervals, lost); err != nil { + t.Fatalf("reconstruct: %v", err) + } + if !bytes.Equal(shardIntervals[lost], want) { + t.Fatalf("rebuilt shard %d does not match the encoded bytes", lost) + } + for _, sid := range spare { + if shardIntervals[sid] != nil { + t.Errorf("shard %d was rebuilt too, only shard %d was asked for", sid, lost) + } + } +}