mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 16:57:45 +02:00
* fix(master): include GrpcPort in LookupEcVolume response LookupVolume already passes loc.GrpcPort through to the client; LookupEcVolume builds Location with only Url / PublicUrl / DataCenter, so callers fall back to ServerToGrpcAddress (httpPort + 10000). On any deployment where that convention does not hold — multi-disk integration tests, custom port layouts — EC reads dial the wrong port and quietly degrade to parity recovery. * fix(volume/ec): probe every DiskLocation when serving local shard reads reconcileEcShardsAcrossDisks (issue 9212) registers each .ec?? against the DiskLocation that physically owns it, so a multi-disk volume server can hold shards for the same vid in two separate ecVolumes — one per disk — with .ecx on whichever disk owned the original .dat. The read path only consulted the single EcVolume FindEcVolume picked, so requests for shards on the sibling disk fell through to errShardNotLocal and then to remote/loopback recovery. Walk all DiskLocations after the first probe in both readLocalEcShardInterval and the VolumeEcShardRead gRPC handler; the latter also covers the loopback that recoverOneRemoteEcShardInterval falls back to when a peer dial fails. * test(volume/ec): cover the multi-disk EC lifecycle end-to-end Two integration tests against a real volume server with two data dirs: TestEcLifecycleAcrossMultipleDisks drives encode -> mount -> HTTP read -> drop .dat -> stop -> redistribute shards across disks -> restart -> verify reconcileEcShardsAcrossDisks attached the orphan shards and reads still work -> blob delete -> stop -> drop a shard -> restart -> VolumeEcShardsRebuild pulls input from both disks -> reads still work. TestEcPartialShardsOnSiblingDiskCleanedUpOnRestart is the issue 9478 reproducer at the cluster level: seed a healthy .dat on disk 0, plant the on-disk footprint of an interrupted EC encode on disk 1, restart, and assert pruneIncompleteEcWithSiblingDat wipes disk 1 without touching disk 0. Framework gets RestartVolumeServer / StopVolumeServer helpers; the previous run's volume.log is rotated to volume.log.previous so a startup regression on the second run does not lose the first run's diagnostics. * review: trim verbose comments * review: drop racy fast-path, use locked findEcShard directly gemini-code-assist flagged the two-step lookup in readLocalEcShardInterval and VolumeEcShardRead: the first probe (ecVolume.FindEcVolumeShard) reads the EcVolume's Shards slice without holding ecVolumesLock, so a concurrent mount / unmount could race with it. findEcShard already walks every DiskLocation under the right lock, so the fast-path adds nothing but the race. Collapse both call sites to a single locked call. Also note in RestartVolumeServer why the log-rotation error is swallowed: absence on first call is benign; anything else surfaces in the next os.Create in startVolume.
590 lines
21 KiB
Go
590 lines
21 KiB
Go
package storage
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"slices"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/klauspost/reedsolomon"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/operation"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/stats"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
|
)
|
|
|
|
// errShardNotLocal indicates that the requested EC shard is simply not
|
|
// stored on this volume server. It is expected during normal reads when
|
|
// shards are spread across multiple servers, so callers should not log
|
|
// it as an error.
|
|
var errShardNotLocal = errors.New("ec shard not on this server")
|
|
|
|
// FindEcShardTargetLocation returns the disk that should receive a new
|
|
// shard / index file for (collection, vid). The selection order is:
|
|
//
|
|
// 1. a disk that already has the EC volume mounted (in-memory state),
|
|
// 2. a disk that owns the .ecx file on disk (volume not mounted yet),
|
|
// 3. any HDD with free space,
|
|
// 4. any disk with free space.
|
|
//
|
|
// Step 2 is the missing primitive that pinned subsequent shards to the
|
|
// first-shard disk during ec.rebuild. ec.rebuild only sets CopyEcxFile=true
|
|
// for the first shard, then relies on auto-select to land later shards on
|
|
// the same disk. Without an on-disk check, FindEcVolume returns nothing
|
|
// (no mount yet) and the fallback picks "any HDD with free space" — which
|
|
// can split shards from their index files across disks of the same node
|
|
// and lose them at startup. See issue #9212 and the orphan-shard
|
|
// reconciliation in #9244.
|
|
//
|
|
// dataShardCount is the data-shard count for this volume's EC layout (10
|
|
// for the OSS default, but custom ratios are supported via .vif). Callers
|
|
// pass it explicitly so this helper stays free of package-level constants
|
|
// — easier to mirror into builds that ship a different default ratio.
|
|
//
|
|
// Implementation walks s.Locations once and scores each disk by tier; the
|
|
// highest-tier disk wins, ties broken by free count. The earlier waterfall
|
|
// across four FindFreeLocation passes was equivalent but acquired
|
|
// volumesLock and ecVolumesLock RLocks (via VolumesLen / EcShardCount) up
|
|
// to four times per disk per call.
|
|
func (s *Store) FindEcShardTargetLocation(collection string, vid needle.VolumeId, dataShardCount int) *DiskLocation {
|
|
const (
|
|
tierAnyDisk = iota + 1
|
|
tierHDD
|
|
tierEcxOnDisk
|
|
tierMounted
|
|
)
|
|
|
|
var (
|
|
best *DiskLocation
|
|
bestTier int
|
|
bestFree int32
|
|
)
|
|
for _, loc := range s.Locations {
|
|
if loc.isDiskSpaceLow {
|
|
continue
|
|
}
|
|
freeCount := ecFreeShardCount(loc, dataShardCount)
|
|
if freeCount <= 0 {
|
|
continue
|
|
}
|
|
tier := tierAnyDisk
|
|
if loc.DiskType == types.HardDriveType {
|
|
tier = tierHDD
|
|
}
|
|
if loc.HasEcxFileOnDisk(collection, vid) {
|
|
tier = tierEcxOnDisk
|
|
}
|
|
if _, mounted := loc.FindEcVolume(vid); mounted {
|
|
tier = tierMounted
|
|
}
|
|
if best == nil || tier > bestTier || (tier == bestTier && freeCount > bestFree) {
|
|
best = loc
|
|
bestTier = tier
|
|
bestFree = freeCount
|
|
}
|
|
}
|
|
return best
|
|
}
|
|
|
|
// ecFreeShardCount returns the free EC shard capacity of loc, expressed
|
|
// in shard slots (not volume-equivalent slots). dataShardCount is the
|
|
// data-shard count of the EC layout being placed — see
|
|
// FindEcShardTargetLocation's docstring for why it's a parameter.
|
|
//
|
|
// FindFreeLocation in store.go does the same math but divides by
|
|
// DataShardsCount at the end. That truncation can exclude a disk that
|
|
// still has room for several individual shards (e.g. MaxVolumeCount=1,
|
|
// EcShardCount=1, dataShardCount=10 → reports 0 despite 9 free shard
|
|
// slots), which in this helper would re-route subsequent shards off the
|
|
// .ecx-owning disk and re-introduce the orphan-shard layout #9212 is
|
|
// trying to prevent. So we keep the result in shard slots throughout.
|
|
//
|
|
// MaxVolumeCount == 0 is the "unlimited" sentinel used elsewhere in the
|
|
// store (see hasFreeDiskLocation). Reporting a synthetic large free
|
|
// count keeps unlimited disks eligible while still letting tie-breaks
|
|
// prefer the less-loaded one.
|
|
func ecFreeShardCount(loc *DiskLocation, dataShardCount int) int32 {
|
|
if dataShardCount <= 0 {
|
|
return 0
|
|
}
|
|
if loc.MaxVolumeCount <= 0 {
|
|
const unlimitedFree = int32(1 << 30)
|
|
used := int32(loc.VolumesLen())*int32(dataShardCount) + int32(loc.EcShardCount())
|
|
if used >= unlimitedFree {
|
|
return 1
|
|
}
|
|
return unlimitedFree - used
|
|
}
|
|
free := (loc.MaxVolumeCount - int32(loc.VolumesLen())) * int32(dataShardCount)
|
|
free -= int32(loc.EcShardCount())
|
|
if free < 0 {
|
|
return 0
|
|
}
|
|
return free
|
|
}
|
|
|
|
func (s *Store) CollectErasureCodingHeartbeat() *master_pb.Heartbeat {
|
|
var ecShardMessages []*master_pb.VolumeEcShardInformationMessage
|
|
collectionEcShardSize := make(map[string]int64)
|
|
for diskId, location := range s.Locations {
|
|
location.ecVolumesLock.RLock()
|
|
for _, ecShards := range location.ecVolumes {
|
|
ecShardMessages = append(ecShardMessages, ecShards.ToVolumeEcShardInformationMessage(uint32(diskId))...)
|
|
|
|
for _, ecShard := range ecShards.Shards {
|
|
collectionEcShardSize[ecShards.Collection] += ecShard.Size()
|
|
}
|
|
}
|
|
location.ecVolumesLock.RUnlock()
|
|
}
|
|
|
|
for col, size := range collectionEcShardSize {
|
|
stats.VolumeServerDiskSizeGauge.WithLabelValues(col, "ec").Set(float64(size))
|
|
}
|
|
|
|
return &master_pb.Heartbeat{
|
|
EcShards: ecShardMessages,
|
|
HasNoEcShards: len(ecShardMessages) == 0,
|
|
}
|
|
|
|
}
|
|
|
|
func (s *Store) MountEcShards(collection string, vid needle.VolumeId, shardId erasure_coding.ShardId, sourceDiskType string) error {
|
|
for diskId, location := range s.Locations {
|
|
if ecVolume, err := location.LoadEcShard(collection, vid, shardId); err == nil {
|
|
glog.V(0).Infof("MountEcShards %d.%d on disk ID %d", vid, shardId, diskId)
|
|
|
|
// Apply the orchestrator-supplied source disk type so the EC
|
|
// volume reports under it instead of the location's. Empty means
|
|
// "fall back to location's disk type" (#9423).
|
|
if sourceDiskType != "" {
|
|
ecVolume.SetDiskType(types.ToDiskType(sourceDiskType))
|
|
}
|
|
|
|
si := erasure_coding.NewShardsInfo()
|
|
si.Set(erasure_coding.NewShardInfo(shardId, erasure_coding.ShardSize(ecVolume.ShardSize())))
|
|
s.NewEcShardsChan <- master_pb.VolumeEcShardInformationMessage{
|
|
Id: uint32(vid),
|
|
Collection: collection,
|
|
EcIndexBits: uint32(si.Bitmap()),
|
|
ShardSizes: si.SizesInt64(),
|
|
DiskType: string(ecVolume.DiskType()),
|
|
ExpireAtSec: ecVolume.ExpireAtSec,
|
|
DiskId: uint32(diskId),
|
|
}
|
|
return nil
|
|
} else if err == os.ErrNotExist {
|
|
continue
|
|
} else {
|
|
return fmt.Errorf("%s load ec shard %d.%d: %v", location.Directory, vid, shardId, err)
|
|
}
|
|
}
|
|
|
|
return fmt.Errorf("MountEcShards %d.%d not found on disk", vid, shardId)
|
|
}
|
|
|
|
func (s *Store) UnmountEcShards(vid needle.VolumeId, shardId erasure_coding.ShardId) error {
|
|
diskId, ecShard, found := s.findEcShard(vid, shardId)
|
|
if !found {
|
|
return nil
|
|
}
|
|
|
|
si := erasure_coding.NewShardsInfo()
|
|
si.Set(erasure_coding.NewShardInfo(shardId, 0))
|
|
message := master_pb.VolumeEcShardInformationMessage{
|
|
Id: uint32(vid),
|
|
Collection: ecShard.Collection,
|
|
EcIndexBits: si.Bitmap(),
|
|
ShardSizes: si.SizesInt64(),
|
|
DiskType: string(ecShard.DiskType),
|
|
DiskId: diskId,
|
|
}
|
|
|
|
location := s.Locations[diskId]
|
|
|
|
if deleted := location.UnloadEcShard(vid, shardId); deleted {
|
|
glog.V(0).Infof("UnmountEcShards %d.%d", vid, shardId)
|
|
s.DeletedEcShardsChan <- message
|
|
return nil
|
|
}
|
|
|
|
return fmt.Errorf("UnmountEcShards %d.%d not found on disk", vid, shardId)
|
|
}
|
|
|
|
func (s *Store) findEcShard(vid needle.VolumeId, shardId erasure_coding.ShardId) (diskId uint32, shard *erasure_coding.EcVolumeShard, found bool) {
|
|
for diskId, location := range s.Locations {
|
|
if v, found := location.FindEcShard(vid, shardId); found {
|
|
return uint32(diskId), v, found
|
|
}
|
|
}
|
|
return 0, nil, false
|
|
}
|
|
|
|
// FindEcShard returns the shard if any DiskLocation on this server holds it,
|
|
// along with that disk's id.
|
|
func (s *Store) FindEcShard(vid needle.VolumeId, shardId erasure_coding.ShardId) (diskId uint32, shard *erasure_coding.EcVolumeShard, found bool) {
|
|
return s.findEcShard(vid, shardId)
|
|
}
|
|
|
|
func (s *Store) FindEcVolume(vid needle.VolumeId) (*erasure_coding.EcVolume, bool) {
|
|
for _, location := range s.Locations {
|
|
if s, found := location.FindEcVolume(vid); found {
|
|
return s, true
|
|
}
|
|
}
|
|
return nil, false
|
|
}
|
|
|
|
// shardFiles is a list of shard files, which is used to return the shard locations
|
|
func (s *Store) CollectEcShards(vid needle.VolumeId, shardFileNames []string) (ecVolume *erasure_coding.EcVolume, found bool) {
|
|
for _, location := range s.Locations {
|
|
if s, foundShards := location.CollectEcShards(vid, shardFileNames); foundShards {
|
|
ecVolume = s
|
|
found = true
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
func (s *Store) DestroyEcVolume(vid needle.VolumeId) {
|
|
for _, location := range s.Locations {
|
|
location.DestroyEcVolume(vid)
|
|
}
|
|
}
|
|
|
|
func (s *Store) ReadEcShardNeedle(vid needle.VolumeId, n *needle.Needle, onReadSizeFn func(size types.Size)) (int, error) {
|
|
for _, location := range s.Locations {
|
|
if localEcVolume, found := location.FindEcVolume(vid); found {
|
|
|
|
offset, size, intervals, err := localEcVolume.LocateEcShardNeedle(n.Id, localEcVolume.Version)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("locate in local ec volume: %w", err)
|
|
}
|
|
if size.IsDeleted() {
|
|
return 0, ErrorDeleted
|
|
}
|
|
|
|
if onReadSizeFn != nil {
|
|
onReadSizeFn(size)
|
|
}
|
|
|
|
glog.V(3).Infof("read ec volume %d offset %d size %d intervals:%+v", vid, offset.ToActualOffset(), size, intervals)
|
|
|
|
if len(intervals) > 1 {
|
|
glog.V(3).Infof("ReadEcShardNeedle needle id %s intervals:%+v", n.String(), intervals)
|
|
}
|
|
bytes, isDeleted, err := s.readEcShardIntervals(n.Id, localEcVolume, intervals)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("ReadEcShardIntervals: %w", err)
|
|
}
|
|
if isDeleted {
|
|
return 0, ErrorDeleted
|
|
}
|
|
|
|
err = n.ReadBytes(bytes, offset.ToActualOffset(), size, localEcVolume.Version)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("ec volume %d needle %s offset %d size %d: %w", vid, n.String(), offset.ToActualOffset(), size, err)
|
|
}
|
|
|
|
return len(bytes), nil
|
|
}
|
|
}
|
|
return 0, fmt.Errorf("ec shard %d not found", vid)
|
|
}
|
|
|
|
func (s *Store) IntervalToShardIdAndOffset(iv erasure_coding.Interval) (erasure_coding.ShardId, int64) {
|
|
return iv.ToShardIdAndOffset(erasure_coding.ErasureCodingLargeBlockSize, erasure_coding.ErasureCodingSmallBlockSize)
|
|
}
|
|
|
|
func (s *Store) readEcShardIntervals(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, intervals []erasure_coding.Interval) (data []byte, is_deleted bool, err error) {
|
|
if err = s.cachedLookupEcShardLocations(ecVolume); err != nil {
|
|
return nil, false, fmt.Errorf("failed to locate shard via master grpc %s: %v", s.MasterAddress, err)
|
|
}
|
|
|
|
for i, interval := range intervals {
|
|
if d, isDeleted, e := s.readOneEcShardInterval(needleId, ecVolume, interval); e != nil {
|
|
return nil, isDeleted, e
|
|
} else {
|
|
if isDeleted {
|
|
is_deleted = true
|
|
}
|
|
if i == 0 {
|
|
data = d
|
|
} else {
|
|
data = append(data, d...)
|
|
}
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
func (s *Store) readOneEcShardInterval(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, interval erasure_coding.Interval) (data []byte, is_deleted bool, err error) {
|
|
shardId, actualOffset := s.IntervalToShardIdAndOffset(interval)
|
|
data = make([]byte, interval.Size)
|
|
|
|
// try local read
|
|
err = s.readLocalEcShardInterval(ecVolume, shardId, data, actualOffset)
|
|
if err == nil {
|
|
return
|
|
}
|
|
if errors.Is(err, errShardNotLocal) {
|
|
// expected when shards are spread across servers; fall through to remote read
|
|
glog.V(4).Infof("ec shard %d.%d not local, will try remote", ecVolume.VolumeId, shardId)
|
|
} else {
|
|
glog.V(0).Infof("read local ec shard %d.%d offset %d: %v", ecVolume.VolumeId, shardId, actualOffset, err)
|
|
}
|
|
|
|
ecVolume.ShardLocationsLock.RLock()
|
|
sourceDataNodes, hasShardIdLocation := ecVolume.ShardLocations[shardId]
|
|
ecVolume.ShardLocationsLock.RUnlock()
|
|
|
|
// try reading directly
|
|
if hasShardIdLocation {
|
|
_, is_deleted, err = s.readRemoteEcShardInterval(sourceDataNodes, needleId, ecVolume.VolumeId, shardId, data, actualOffset)
|
|
if err == nil {
|
|
return
|
|
}
|
|
glog.V(0).Infof("read remote ec shard %d.%d locations: %v", ecVolume.VolumeId, shardId, err)
|
|
}
|
|
|
|
// try reading by recovering from other shards
|
|
_, is_deleted, err = s.recoverOneRemoteEcShardInterval(needleId, ecVolume, shardId, data, actualOffset)
|
|
if err == nil {
|
|
return
|
|
}
|
|
glog.V(0).Infof("recover ec shard %d.%d : %v", ecVolume.VolumeId, shardId, err)
|
|
|
|
return
|
|
}
|
|
|
|
func forgetShardId(ecVolume *erasure_coding.EcVolume, shardId erasure_coding.ShardId) {
|
|
// failed to access the source data nodes, clear it up
|
|
ecVolume.ShardLocationsLock.Lock()
|
|
delete(ecVolume.ShardLocations, shardId)
|
|
ecVolume.ShardLocationsLock.Unlock()
|
|
}
|
|
|
|
func (s *Store) cachedLookupEcShardLocations(ecVolume *erasure_coding.EcVolume) (err error) {
|
|
|
|
shardCount := len(ecVolume.ShardLocations)
|
|
if shardCount < erasure_coding.DataShardsCount &&
|
|
ecVolume.ShardLocationsRefreshTime.Add(11*time.Second).After(time.Now()) ||
|
|
shardCount == erasure_coding.TotalShardsCount &&
|
|
ecVolume.ShardLocationsRefreshTime.Add(37*time.Minute).After(time.Now()) ||
|
|
shardCount >= erasure_coding.DataShardsCount &&
|
|
ecVolume.ShardLocationsRefreshTime.Add(7*time.Minute).After(time.Now()) {
|
|
// still fresh
|
|
return nil
|
|
}
|
|
|
|
glog.V(3).Infof("lookup and cache ec volume %d locations", ecVolume.VolumeId)
|
|
|
|
err = operation.WithMasterServerClient(false, s.MasterAddress, s.grpcDialOption, func(masterClient master_pb.SeaweedClient) error {
|
|
req := &master_pb.LookupEcVolumeRequest{
|
|
VolumeId: uint32(ecVolume.VolumeId),
|
|
}
|
|
resp, err := masterClient.LookupEcVolume(context.Background(), req)
|
|
if err != nil {
|
|
return fmt.Errorf("lookup ec volume %d: %v", ecVolume.VolumeId, err)
|
|
}
|
|
if len(resp.ShardIdLocations) < erasure_coding.DataShardsCount {
|
|
return fmt.Errorf("only %d shards found but %d required", len(resp.ShardIdLocations), erasure_coding.DataShardsCount)
|
|
}
|
|
|
|
ecVolume.ShardLocationsLock.Lock()
|
|
for _, shardIdLocations := range resp.ShardIdLocations {
|
|
shardId := erasure_coding.ShardId(shardIdLocations.ShardId)
|
|
delete(ecVolume.ShardLocations, shardId)
|
|
for _, loc := range shardIdLocations.Locations {
|
|
ecVolume.ShardLocations[shardId] = append(ecVolume.ShardLocations[shardId], pb.NewServerAddressFromLocation(loc))
|
|
}
|
|
}
|
|
ecVolume.ShardLocationsRefreshTime = time.Now()
|
|
ecVolume.ShardLocationsLock.Unlock()
|
|
|
|
return nil
|
|
})
|
|
return
|
|
}
|
|
|
|
func (s *Store) readLocalEcShardInterval(ecVolume *erasure_coding.EcVolume, shardId erasure_coding.ShardId, buf []byte, offset int64) error {
|
|
// findEcShard walks every DiskLocation under ecVolumesLock; the
|
|
// shard may live on a sibling disk of this server.
|
|
_, shard, found := s.findEcShard(ecVolume.VolumeId, shardId)
|
|
if !found {
|
|
return fmt.Errorf("shard %d for volume %d: %w", shardId, ecVolume.VolumeId, errShardNotLocal)
|
|
}
|
|
|
|
readBytes, err := shard.ReadAt(buf, offset)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to read local EC shard %d for volume %d: %v", shardId, ecVolume.VolumeId, err)
|
|
}
|
|
if got, want := readBytes, len(buf); got != want {
|
|
return fmt.Errorf("expected %d bytes for local EC shard %d on volume %d, got %d", want, shardId, ecVolume.VolumeId, got)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) readRemoteEcShardInterval(sourceDataNodes []pb.ServerAddress, needleId types.NeedleId, vid needle.VolumeId, shardId erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) {
|
|
|
|
if len(sourceDataNodes) == 0 {
|
|
return 0, false, fmt.Errorf("failed to find ec shard %d.%d", vid, shardId)
|
|
}
|
|
|
|
for _, sourceDataNode := range sourceDataNodes {
|
|
glog.V(3).Infof("read remote ec shard %d.%d from %s", vid, shardId, sourceDataNode)
|
|
n, is_deleted, err = s.doReadRemoteEcShardInterval(sourceDataNode, needleId, vid, shardId, buf, offset)
|
|
if err == nil {
|
|
return
|
|
}
|
|
glog.V(1).Infof("read remote ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
func (s *Store) doReadRemoteEcShardInterval(sourceDataNode pb.ServerAddress, needleId types.NeedleId, vid needle.VolumeId, shardId erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) {
|
|
|
|
err = operation.WithVolumeServerClient(false, sourceDataNode, s.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
|
|
|
|
// copy data slice
|
|
shardReadClient, err := client.VolumeEcShardRead(context.Background(), &volume_server_pb.VolumeEcShardReadRequest{
|
|
VolumeId: uint32(vid),
|
|
ShardId: uint32(shardId),
|
|
Offset: offset,
|
|
Size: int64(len(buf)),
|
|
FileKey: uint64(needleId),
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("failed to start reading ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
|
|
}
|
|
|
|
for {
|
|
resp, receiveErr := shardReadClient.Recv()
|
|
if receiveErr == io.EOF {
|
|
break
|
|
}
|
|
if receiveErr != nil {
|
|
return fmt.Errorf("receiving ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, receiveErr)
|
|
}
|
|
if resp.IsDeleted {
|
|
is_deleted = true
|
|
}
|
|
copy(buf[n:n+len(resp.Data)], resp.Data)
|
|
n += len(resp.Data)
|
|
}
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return 0, is_deleted, fmt.Errorf("read ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
|
|
}
|
|
|
|
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)
|
|
|
|
enc, err := reedsolomon.New(erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
|
|
if err != nil {
|
|
return 0, false, fmt.Errorf("failed to create encoder: %w", err)
|
|
}
|
|
|
|
// Use MaxShardCount to support custom EC ratios up to 32 shards
|
|
bufs := make([][]byte, erasure_coding.MaxShardCount)
|
|
|
|
var wg sync.WaitGroup
|
|
ecVolume.ShardLocationsLock.RLock()
|
|
for shardId, locations := range ecVolume.ShardLocations {
|
|
|
|
// skip current shard or empty shard
|
|
if shardId == shardIdToRecover {
|
|
continue
|
|
}
|
|
if len(locations) == 0 {
|
|
glog.V(3).Infof("readRemoteEcShardInterval missing %d.%d from %+v", ecVolume.VolumeId, shardId, locations)
|
|
continue
|
|
}
|
|
|
|
// read from remote locations
|
|
wg.Add(1)
|
|
go func(shardId erasure_coding.ShardId, locations []pb.ServerAddress) {
|
|
defer wg.Done()
|
|
data := make([]byte, len(buf))
|
|
nRead, isDeleted, readErr := s.readRemoteEcShardInterval(locations, needleId, ecVolume.VolumeId, shardId, data, offset)
|
|
if readErr != nil {
|
|
glog.V(3).Infof("recover: readRemoteEcShardInterval %d.%d %d bytes from %+v: %v", ecVolume.VolumeId, shardId, nRead, locations, readErr)
|
|
forgetShardId(ecVolume, shardId)
|
|
}
|
|
if isDeleted {
|
|
is_deleted = true
|
|
}
|
|
if nRead == len(buf) {
|
|
bufs[shardId] = data
|
|
}
|
|
}(shardId, locations)
|
|
}
|
|
ecVolume.ShardLocationsLock.RUnlock()
|
|
|
|
wg.Wait()
|
|
|
|
// Count and log available shards for diagnostics
|
|
availableShards := make([]erasure_coding.ShardId, 0, erasure_coding.TotalShardsCount)
|
|
missingShards := make([]erasure_coding.ShardId, 0, erasure_coding.ParityShardsCount+1)
|
|
for shardId := 0; shardId < erasure_coding.TotalShardsCount; shardId++ {
|
|
if bufs[shardId] != nil {
|
|
availableShards = append(availableShards, erasure_coding.ShardId(shardId))
|
|
} else {
|
|
missingShards = append(missingShards, erasure_coding.ShardId(shardId))
|
|
}
|
|
}
|
|
|
|
glog.V(3).Infof("recover ec shard %d.%d: %d shards available %v, %d missing %v",
|
|
ecVolume.VolumeId, shardIdToRecover,
|
|
len(availableShards), availableShards,
|
|
len(missingShards), missingShards)
|
|
|
|
if len(availableShards) < erasure_coding.DataShardsCount {
|
|
return 0, false, fmt.Errorf("cannot recover shard %d.%d: only %d shards available %v, need at least %d (missing: %v)",
|
|
ecVolume.VolumeId, shardIdToRecover,
|
|
len(availableShards), availableShards,
|
|
erasure_coding.DataShardsCount, missingShards)
|
|
}
|
|
|
|
if err = enc.ReconstructData(bufs[:erasure_coding.TotalShardsCount]); err != nil {
|
|
return 0, false, fmt.Errorf("failed to reconstruct data for shard %d.%d with %d available shards %v: %w",
|
|
ecVolume.VolumeId, shardIdToRecover, len(availableShards), availableShards, err)
|
|
}
|
|
glog.V(4).Infof("recovered ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover)
|
|
|
|
copy(buf, bufs[shardIdToRecover])
|
|
|
|
return len(buf), is_deleted, nil
|
|
}
|
|
|
|
func (s *Store) EcVolumes() (ecVolumes []*erasure_coding.EcVolume) {
|
|
for _, location := range s.Locations {
|
|
location.ecVolumesLock.RLock()
|
|
for _, v := range location.ecVolumes {
|
|
ecVolumes = append(ecVolumes, v)
|
|
}
|
|
location.ecVolumesLock.RUnlock()
|
|
}
|
|
slices.SortFunc(ecVolumes, func(a, b *erasure_coding.EcVolume) int {
|
|
return int(a.VolumeId) - int(b.VolumeId)
|
|
})
|
|
return ecVolumes
|
|
}
|