package topology import ( "context" "fmt" "math/rand/v2" "sync" "sync/atomic" "time" "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" "github.com/seaweedfs/seaweedfs/weed/storage/types" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/storage" "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/storage/super_block" ) type copyState int const ( noCopies copyState = 0 + iota insufficientCopies enoughCopies ) type volumeState string const ( readOnlyState volumeState = "ReadOnly" oversizedState = "Oversized" crowdedState = "Crowded" NoWritableVolumes = "No writable volumes" ) // volumeSizeTracking holds per-volume size accounting for weighted assignment. type volumeSizeTracking struct { effectiveSize uint64 // reported + pending assigned bytes reportedSize uint64 // last heartbeat-reported size (dedup replicas) compactRevision uint32 // detect compaction to reset instead of decay lastUpdateTime time.Time // dedup replicas within the same heartbeat cycle fullSince time.Time // non-zero while the volume is marked full (RecordAssign or capacity heartbeat) } // capacityRecoveryDelay is the minimum time a volume must stay out of the // writable list after being removed for capacity before it can // be considered for re-addition by heartbeat-driven decay. Combined with // the effectiveSize hysteresis band, this avoids bouncing the volume in // and out of writable within a single burst of assigns. const capacityRecoveryDelay = 30 * time.Second // mapping from volume to its locations, inverted from server to volume type VolumeLayout struct { growRequest atomic.Bool lastGrowCount atomic.Uint32 rp *super_block.ReplicaPlacement ttl *needle.TTL diskType types.DiskType vid2location map[needle.VolumeId]*VolumeLocationList writables []needle.VolumeId // transient array of writable volume id crowded map[needle.VolumeId]struct{} vacuumedVolumes map[needle.VolumeId]time.Time volumeSizeLimit uint64 replicationAsMin bool accessLock sync.RWMutex sizeTracking map[needle.VolumeId]*volumeSizeTracking // dropped: the layout went away with its collection; late registrations must re-resolve. dropped bool } type VolumeLayoutStats struct { TotalSize uint64 UsedSize uint64 // LogicalUsedSize counts one copy of the data: a single replica of a // regular volume, the data shards of an EC volume. LogicalUsedSize uint64 FileCount uint64 } func NewVolumeLayout(rp *super_block.ReplicaPlacement, ttl *needle.TTL, diskType types.DiskType, volumeSizeLimit uint64, replicationAsMin bool) *VolumeLayout { return &VolumeLayout{ rp: rp, ttl: ttl, diskType: diskType, vid2location: make(map[needle.VolumeId]*VolumeLocationList), writables: *new([]needle.VolumeId), crowded: make(map[needle.VolumeId]struct{}), vacuumedVolumes: make(map[needle.VolumeId]time.Time), volumeSizeLimit: volumeSizeLimit, replicationAsMin: replicationAsMin, sizeTracking: make(map[needle.VolumeId]*volumeSizeTracking), } } func (vl *VolumeLayout) String() string { return fmt.Sprintf("rp:%v, ttl:%v, writables:%v, volumeSizeLimit:%v", vl.rp, vl.ttl, vl.writables, vl.volumeSizeLimit) } // getOrCreateLocationList returns the vid's location list, creating an empty one // if the index has none. Callers hold accessLock. func (vl *VolumeLayout) getOrCreateLocationList(vid needle.VolumeId) *VolumeLocationList { location, ok := vl.vid2location[vid] if !ok { location = NewVolumeLocationList() vl.vid2location[vid] = location } return location } // initSizeTracking seeds a vid's size tracking from a reported size if it has // none yet. Callers hold accessLock. func (vl *VolumeLayout) initSizeTracking(vid needle.VolumeId, size uint64, compactRevision uint32) { if _, exists := vl.sizeTracking[vid]; !exists { vl.sizeTracking[vid] = &volumeSizeTracking{ effectiveSize: size, reportedSize: size, compactRevision: compactRevision, } } } // RegisterVolume records a volume location. It refuses a layout dropped with // its collection — the caller must re-resolve, or the bits it sets would leak. func (vl *VolumeLayout) RegisterVolume(v *storage.VolumeInfo, dn *DataNode) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock() if vl.dropped { return false } moveLookupOwnership(v.Id, vl.getOrCreateLocationList(v.Id).Set(dn), dn) if !v.ReadOnly { vl.initSizeTracking(v.Id, v.Size, v.CompactRevision) } // glog.V(4).Infof("volume %d added to %s len %d copy %d", v.Id, dn.Id(), vl.vid2location[v.Id].Length(), v.ReplicaPlacement.GetCopyCount()) location := vl.vid2location[v.Id] for _, dn := range location.list { if vInfo, err := dn.GetVolumesById(v.Id); err == nil { if vInfo.ReadOnly { glog.V(1).Infof("vid %d removed from writable", v.Id) vl.removeFromWritable(v.Id) location.SetReadOnly(dn, true) return true } location.SetReadOnly(dn, false) } else { glog.V(1).Infof("vid %d removed from writable", v.Id) vl.removeFromWritable(v.Id) location.SetReadOnly(dn, false) return true } } return true } func (vl *VolumeLayout) rememberOversizedVolume(v *storage.VolumeInfo, dn *DataNode) { if location, ok := vl.vid2location[v.Id]; ok { location.SetOversized(dn, vl.isOversized(v)) } } // UpdateOversizedState refreshes the per-replica oversized mark from the size // reported by this heartbeat, and is the only writer of that mark. The delta // heartbeat path (ApplyVolumeChanges) and the full heartbeat path // (SyncDataNodeRegistration) both call it, or a volume that grew past the // limit would keep looking writable to ensureCorrectWritables and bounce // between writable and unwritable on every assign. // // Registration deliberately does not touch the mark: the incremental path // registers from a short heartbeat message that carries no size, so judging // the volume by it would clear the mark on every arrival announcement and // hand an oversized volume back to the writable list. func (vl *VolumeLayout) UpdateOversizedState(v *storage.VolumeInfo, dn *DataNode) { vl.accessLock.Lock() defer vl.accessLock.Unlock() vl.rememberOversizedVolume(v, dn) } // UpdateVolumeSize is called on every heartbeat for every reported volume. // It decays the pending size estimate toward the reported size and updates // crowded state. Replicated volumes report from multiple DataNodes; decay // runs only once per new reported size to avoid double-halving. // If the compact revision changed, the size drop is from compaction (not // pending writes), so we reset effectiveSize to the reported size instead of // decaying. // // Returns recoveredToWritable = true when decay brought a volume that was // previously eagerly removed by RecordAssign back under the writable // threshold and this call re-added it to the writable list. The caller // should mirror the activeVolumeCount bookkeeping. // // fromHeartbeat marks a volume server's report. The periodic decay stands in // for the report a quiet volume never sends, so it passes no size of its own — // it reads the record under this lock, and never advances lastUpdateTime, or // it would swallow the next real report and strand the master on a stale size // no one will send again. Both callers give way to a report already handled // for this cycle. func (vl *VolumeLayout) UpdateVolumeSize(vid needle.VolumeId, reportedSize uint64, compactRevision uint32, fromHeartbeat bool) (recoveredToWritable bool) { vl.accessLock.Lock() defer vl.accessLock.Unlock() // Tracking exists to place writes, and nothing is written to a volume any // replica reports read-only. Most volumes in a tiered cluster are read-only, // so an entry each is the layout's largest cost. Asked of the volume rather // than of the replica reporting, so which replica arrives first cannot // decide the answer. A volume held out of the writable list for capacity is // still all-writable and keeps its entry, which is what enforces the // recovery delay. if !vl.isAllWritable(vid) { delete(vl.sizeTracking, vid) vl.removeFromCrowded(vid) return false } now := time.Now() st := vl.sizeTracking[vid] if st == nil { if !fromHeartbeat { return false // the entry went while the decay was picking its work } st = &volumeSizeTracking{ effectiveSize: reportedSize, reportedSize: reportedSize, compactRevision: compactRevision, lastUpdateTime: now, } vl.sizeTracking[vid] = st } else if now.Sub(st.lastUpdateTime) < 2*time.Second { // Something already decayed this volume for this cycle: another replica // of the same report, or for the decay a heartbeat that arrived after // the pass chose the volume. Halving twice would forget pending bytes // the volume has not written yet. return false } else { if fromHeartbeat { st.lastUpdateTime = now st.reportedSize = reportedSize } else { // Take the record as it stands under this lock, not as the decay // pass saw it: a heartbeat may have landed since, and replaying // the older size would roll its report back. reportedSize, compactRevision = st.reportedSize, st.compactRevision } if compactRevision != st.compactRevision { // Compaction happened — size drop is real, not pending. Reset. st.compactRevision = compactRevision st.effectiveSize = reportedSize } else if st.effectiveSize > reportedSize { st.effectiveSize = reportedSize + (st.effectiveSize-reportedSize)/2 } else { st.effectiveSize = reportedSize } } crowdedThreshold := uint64(float64(vl.volumeSizeLimit) * VolumeGrowStrategy.Threshold) if st.effectiveSize > crowdedThreshold { vl.setVolumeCrowded(vid) return false } vl.removeFromCrowded(vid) // Recovery path: if we eagerly removed this volume from writables in // RecordAssign, decay may now have brought effectiveSize back under // the crowded threshold. Re-add it — but only after the recovery // delay has elapsed, so a steady stream of assigns near the limit // does not bounce the volume in and out of writables. if st.fullSince.IsZero() { return false } if now.Sub(st.fullSince) < capacityRecoveryDelay { return false } if reportedSize >= vl.volumeSizeLimit { return false // actual on-disk size still over limit; stay out } if vl.vid2location[vid].AnyOversized() { return false } if !vl.enoughCopies(vid) || !vl.isAllWritable(vid) { return false } if !vl.setVolumeWritable(vid) { return false // already writable (shouldn't happen, but be safe) } st.fullSince = time.Time{} glog.V(0).Infof("Volume %d recovered to writable (effective=%d, reported=%d, limit=%d).", vid, st.effectiveSize, reportedSize, vl.volumeSizeLimit) return true } // DecayQuietVolumeSizes re-runs the size decay for volumes whose heartbeat // reports have gone quiet. Only a volume whose content changed is reported, // and a volume held out of the writable list takes no writes, so without this // an inflated estimate is never decayed and the volume never returns. // // A volume the disk really did fill is skipped: UpdateVolumeSize refuses to // recover one whose reported size is at the limit, so walking it every pulse // only takes the write lock away from the heartbeats. Nothing is lost by // waiting — shrinking it means compaction, which changes content and is // therefore reported. func (vl *VolumeLayout) DecayQuietVolumeSizes(quietCutoff time.Duration) { for _, vid := range vl.quietDecayCandidates(quietCutoff) { if vl.UpdateVolumeSize(vid, 0, 0, false) { vl.AdjustActiveVolumeCountAfterRecovery(vid) } } } // quietDecayCandidates names the volumes the decay has something to do for. // Skipping the rest is what keeps the pass off the write lock: the work is // what costs, not the outcome, since replaying a size that cannot move leaves // the record exactly as it found it. func (vl *VolumeLayout) quietDecayCandidates(quietCutoff time.Duration) (quiets []needle.VolumeId) { now := time.Now() vl.accessLock.RLock() defer vl.accessLock.RUnlock() for vid, st := range vl.sizeTracking { if now.Sub(st.lastUpdateTime) < quietCutoff { continue } if st.effectiveSize > st.reportedSize || (!st.fullSince.IsZero() && st.reportedSize < vl.volumeSizeLimit) { quiets = append(quiets, vid) } } return quiets } func (vl *VolumeLayout) UnRegisterVolume(v *storage.VolumeInfo, dn *DataNode) { vl.accessLock.Lock() defer vl.accessLock.Unlock() // remove from vid2location map location, ok := vl.vid2location[v.Id] if !ok { return } if removed := location.Remove(dn); removed != nil { moveLookupOwnership(v.Id, removed, nil) vl.ensureCorrectWritables(v.Id) if location.Length() == 0 { delete(vl.vid2location, v.Id) delete(vl.sizeTracking, v.Id) vl.removeFromCrowded(v.Id) } } } func (vl *VolumeLayout) EnsureCorrectWritables(v *storage.VolumeInfo) { vl.accessLock.Lock() defer vl.accessLock.Unlock() vl.ensureCorrectWritables(v.Id) } func (vl *VolumeLayout) ensureCorrectWritables(vid needle.VolumeId) { isEnoughCopies := vl.enoughCopies(vid) isAllWritable := vl.isAllWritable(vid) isOversizedVolume := vl.vid2location[vid].AnyOversized() if isEnoughCopies && isAllWritable && !isOversizedVolume { // A volume removed for capacity (fullSince set) must go through the // heartbeat recovery path in UpdateVolumeSize, which enforces // capacityRecoveryDelay. Re-adding it here would bypass the cooldown // and let a just-compacted volume flip back to writable immediately. if st := vl.sizeTracking[vid]; st != nil && !st.fullSince.IsZero() && time.Since(st.fullSince) < capacityRecoveryDelay { return } // A volume already at the limit must not be restored here: the next // assign's RecordAssign would remove it again. Crowded is the wrong // test for that -- it starts at 90%, and a crowded volume that lost a // replica would never be writable again once the replica returned, // since nothing writes to it and its size can no longer fall. if st := vl.sizeTracking[vid]; st != nil && st.effectiveSize >= vl.volumeSizeLimit { return } vl.setVolumeWritable(vid) return } // removeFromWritable reports the transition itself, and only when there is // one. Explain it only then: every heartbeat re-runs this for every volume, // so a volume that is simply staying read-only must not log. if !vl.removeFromWritable(vid) { return } if !isEnoughCopies { glog.V(0).Infof("volume %d does not have enough copies", vid) } if !isAllWritable { glog.V(0).Infof("volume %d is not fully writable", vid) } if isOversizedVolume { glog.V(0).Infof("volume %d is oversized", vid) } } func (vl *VolumeLayout) isAllWritable(vid needle.VolumeId) bool { if location, ok := vl.vid2location[vid]; ok { for _, dn := range location.list { if v, getError := dn.GetVolumesById(vid); getError == nil { if v.ReadOnly { return false } } } } else { return false } return true } func (vl *VolumeLayout) isOversized(v *storage.VolumeInfo) bool { return uint64(v.Size) >= vl.volumeSizeLimit } // RecordAssign adds the estimated byte size to the volume's tracked effective // size and updates the volume's writable state: // // - at the crowded threshold (e.g. 90%): marks crowded to trigger growth. // - at the hard limit (100%): removes the volume from the writable list // immediately so new assigns stop landing on it. Returns true if this // call was what removed the volume, so the caller can mirror the // disk-usage accounting done by Topology.SetVolumeCapacityFull. // // Removing eagerly here avoids waiting for the heartbeat-driven // CollectDeadNodeAndFullVolumes cycle (5–15s detection latency) during // which a fast writer could push the volume far past the configured limit. func (vl *VolumeLayout) RecordAssign(vid needle.VolumeId, pendingDelta int64) (reachedCapacity bool) { vl.accessLock.Lock() defer vl.accessLock.Unlock() st := vl.sizeTracking[vid] if st == nil { return false } if pendingDelta > 0 { st.effectiveSize += uint64(pendingDelta) } if st.effectiveSize >= vl.volumeSizeLimit { if vl.removeFromWritable(vid) { st.fullSince = time.Now() glog.V(0).Infof("Volume %d reaches full capacity (effective=%d, limit=%d).", vid, st.effectiveSize, vl.volumeSizeLimit) return true } return false } if float64(st.effectiveSize) > float64(vl.volumeSizeLimit)*VolumeGrowStrategy.Threshold { vl.setVolumeCrowded(vid) } return false } // AdjustActiveVolumeCountForFull decrements the active volume count on each // data node holding this volume. Mirrors the accounting done in // Topology.SetVolumeCapacityFull for the heartbeat-driven path. Call only // after RecordAssign returns true for the same vid. func (vl *VolumeLayout) AdjustActiveVolumeCountForFull(vid needle.VolumeId) { vl.adjustActiveVolumeCount(vid, -1) } // AdjustActiveVolumeCountAfterRecovery increments the active volume count on // each data node holding this volume. Mirrors // AdjustActiveVolumeCountForFull for the recovery path. Call only after // UpdateVolumeSize returns true for the same vid. func (vl *VolumeLayout) AdjustActiveVolumeCountAfterRecovery(vid needle.VolumeId) { vl.adjustActiveVolumeCount(vid, +1) } func (vl *VolumeLayout) adjustActiveVolumeCount(vid needle.VolumeId, delta int64) { // Copy the node list under the VolumeLayout lock, then release it before // calling UpAdjustDiskUsageDelta. UpAdjustDiskUsageDelta walks up the // topology tree taking per-level locks (e.g., DiskUsages.Lock on each // node). Keeping vl.accessLock held across that tree walk is an // unnecessary lock-ordering hazard — other call paths that hold a // topology-level lock and then need vl.accessLock would deadlock. vl.accessLock.RLock() vidLocations, found := vl.vid2location[vid] if !found { vl.accessLock.RUnlock() return } nodes := make([]*DataNode, len(vidLocations.list)) copy(nodes, vidLocations.list) vl.accessLock.RUnlock() diskTypeStr := string(vl.diskType) for _, dn := range nodes { disk := dn.getOrCreateDisk(diskTypeStr) disk.UpAdjustDiskUsageDelta(vl.diskType, &DiskUsageCounts{ activeVolumeCount: delta, }) } } const maxDrainWait = 30 * time.Second const pendingSizeThreshold uint64 = 2 * 1024 * 1024 // 2 MB // GetPendingSize returns the estimated in-flight bytes for a volume: // the gap between the effective tracked size and the last heartbeat-reported size. func (vl *VolumeLayout) GetPendingSize(vid needle.VolumeId) uint64 { vl.accessLock.RLock() defer vl.accessLock.RUnlock() if st := vl.sizeTracking[vid]; st != nil && st.effectiveSize > st.reportedSize { return st.effectiveSize - st.reportedSize } return 0 } // waitForPendingDrain polls until pending bytes for the volume decay below // the threshold, the timeout expires, or the context is cancelled. Since the // volume is already removed from the writable list, no new assigns accumulate // — pending only decreases via heartbeat decay. func (vl *VolumeLayout) waitForPendingDrain(ctx context.Context, vid needle.VolumeId) { deadline := time.Now().Add(maxDrainWait) for time.Now().Before(deadline) { if vl.GetPendingSize(vid) <= pendingSizeThreshold { return } select { case <-ctx.Done(): return case <-time.After(1 * time.Second): } } glog.Warningf("volume %d: %d pending bytes remain after drain timeout", vid, vl.GetPendingSize(vid)) } // DrainAndRemoveFromWritable removes the volume from the writable list // immediately, then waits for pending assigned bytes to decay. // Used by vacuum before compaction. func (vl *VolumeLayout) DrainAndRemoveFromWritable(vid needle.VolumeId) { vl.accessLock.Lock() vl.removeFromWritable(vid) vl.accessLock.Unlock() vl.waitForPendingDrain(context.Background(), vid) } func (vl *VolumeLayout) isEmpty() bool { vl.accessLock.RLock() defer vl.accessLock.RUnlock() return len(vl.vid2location) == 0 } func (vl *VolumeLayout) Lookup(vid needle.VolumeId) []*DataNode { vl.accessLock.RLock() defer vl.accessLock.RUnlock() if location := vl.vid2location[vid]; location != nil { return location.list } return nil } // HasDataNode reports whether the layout already lists dn as a location for vid. // Used to detect a volume that is present on a data node but missing from the // lookup index, so it can be re-registered. func (vl *VolumeLayout) HasDataNode(vid needle.VolumeId, dn *DataNode) bool { vl.accessLock.RLock() defer vl.accessLock.RUnlock() location, ok := vl.vid2location[vid] if !ok { return false } for _, n := range location.list { if n.Ip == dn.Ip && n.Port == dn.Port { return true } } return false } func (vl *VolumeLayout) ListVolumeServers() (nodes []*DataNode) { vl.accessLock.RLock() defer vl.accessLock.RUnlock() for _, location := range vl.vid2location { nodes = append(nodes, location.list...) } return } func (vl *VolumeLayout) PickForWrite(count uint64, option *VolumeGrowOption) (vid needle.VolumeId, counter uint64, locationList *VolumeLocationList, shouldGrow bool, err error) { vl.accessLock.RLock() defer vl.accessLock.RUnlock() lenWriters := len(vl.writables) if lenWriters <= 0 { return 0, 0, nil, true, fmt.Errorf("%s", NoWritableVolumes) } if option.DataCenter == "" && option.Rack == "" && option.DataNode == "" { vid, locationList = vl.pickWeightedByRemaining(vl.writables) if locationList == nil || len(locationList.list) == 0 { return 0, 0, nil, false, fmt.Errorf("Strangely vid %s is on no machine!", vid.String()) } return vid, count, locationList.Copy(), false, nil } // Scan from a random offset to collect up to pickSampleSize matching // candidates, avoiding a full scan + allocation in the common case. var sample [pickSampleSize]needle.VolumeId found := 0 start := rand.IntN(lenWriters) for i := 0; i < lenWriters && found < pickSampleSize; i++ { writableVolumeId := vl.writables[(start+i)%lenWriters] volumeLocationList := vl.vid2location[writableVolumeId] for _, dn := range volumeLocationList.list { if option.DataCenter != "" && dn.GetDataCenter().Id() != NodeId(option.DataCenter) { continue } if option.Rack != "" && dn.GetRack().Id() != NodeId(option.Rack) { continue } if option.DataNode != "" && dn.Id() != NodeId(option.DataNode) { continue } sample[found] = writableVolumeId found++ break } } if found == 0 { return vid, count, locationList, true, fmt.Errorf("%s in DataCenter:%v Rack:%v DataNode:%v", NoWritableVolumes, option.DataCenter, option.Rack, option.DataNode) } vid, locationList = vl.weightedPick(sample[:found]) return vid, count, locationList.Copy(), false, nil } // pickSampleSize is how many random candidates to sample before doing a // weighted pick. Keeps cost O(1) regardless of total writable volume count // while still biasing toward emptier volumes. const pickSampleSize = 3 // pickWeightedByRemaining randomly samples a few candidates from the list, // then does a weighted pick among them by remaining capacity. // Sampled candidates may repeat when len(candidates) is small relative to // pickSampleSize; this is harmless — a repeated volume just gets proportionally // more weight, which is a negligible statistical effect. func (vl *VolumeLayout) pickWeightedByRemaining(candidates []needle.VolumeId) (needle.VolumeId, *VolumeLocationList) { n := len(candidates) if n <= pickSampleSize { return vl.weightedPick(candidates) } var sample [pickSampleSize]needle.VolumeId for i := range sample { sample[i] = candidates[rand.IntN(n)] } return vl.weightedPick(sample[:]) } func (vl *VolumeLayout) weightedPick(candidates []needle.VolumeId) (needle.VolumeId, *VolumeLocationList) { if len(candidates) == 1 { vid := candidates[0] return vid, vl.vid2location[vid] } // first pass: sum weights var totalRemaining uint64 for _, vid := range candidates { totalRemaining += vl.remainingSize(vid) } // second pass: weighted random pick pick := rand.Uint64N(totalRemaining) var cumulative uint64 for _, vid := range candidates { cumulative += vl.remainingSize(vid) if pick < cumulative { return vid, vl.vid2location[vid] } } vid := candidates[0] return vid, vl.vid2location[vid] } func (vl *VolumeLayout) remainingSize(vid needle.VolumeId) uint64 { var size uint64 if st := vl.sizeTracking[vid]; st != nil { size = st.effectiveSize } if size < vl.volumeSizeLimit { if r := vl.volumeSizeLimit - size; r > 1 { return r } } return 1 } func (vl *VolumeLayout) HasGrowRequest() bool { return vl.growRequest.Load() } // AddGrowRequestIfAbsent atomically claims the pending-growth flag. It returns // true for the one caller that transitions it from unset to set (the growth // initiator); concurrent callers get false and are followers of that growth. func (vl *VolumeLayout) AddGrowRequestIfAbsent() bool { return vl.growRequest.CompareAndSwap(false, true) } func (vl *VolumeLayout) DoneGrowRequest() { vl.growRequest.Store(false) } func (vl *VolumeLayout) SetLastGrowCount(count uint32) { if vl.lastGrowCount.Load() != count && count != 0 { vl.lastGrowCount.Store(count) } } func (vl *VolumeLayout) GetLastGrowCount() uint32 { return vl.lastGrowCount.Load() } func (vl *VolumeLayout) ShouldGrowVolumes() bool { writable, crowded := vl.GetWritableVolumeCount() return writable <= crowded } func (vl *VolumeLayout) ShouldGrowVolumesByDcAndRack(writables *[]needle.VolumeId, dcId NodeId, rackId NodeId) bool { // When replication spans multiple racks (DiffRackCount > 0), a writable // volume's replicas only cover some racks in a DC. It is wrong to // require every rack to host a replica — that would create volumes // endlessly in any DC with more racks than the copy count. // Instead, check at the DC level: if the DC already has a non-crowded // writable volume, no growth is needed for uncovered racks. checkDcOnly := vl.rp.DiffRackCount > 0 for _, v := range *writables { for _, dn := range vl.Lookup(v) { if dn.GetDataCenter().Id() != dcId { continue } if !checkDcOnly && dn.GetRack().Id() != rackId { continue } if _, err := dn.GetVolumesById(v); err == nil { vl.accessLock.RLock() var size uint64 if st := vl.sizeTracking[v]; st != nil { size = st.effectiveSize } vl.accessLock.RUnlock() if float64(size) <= float64(vl.volumeSizeLimit)*VolumeGrowStrategy.Threshold { return false } } } } return true } // listVolumeDataCenters returns the data centers hosting at least one volume // of this layout, stopping once limit distinct DCs are seen (0 = no limit). // The walk holds accessLock and a layout can span millions of volumes, so // callers pass the smallest limit that answers their question and only a // layout truly confined to fewer DCs pays for a full walk. func (vl *VolumeLayout) listVolumeDataCenters(limit int) map[NodeId]struct{} { vl.accessLock.RLock() defer vl.accessLock.RUnlock() dataCenters := make(map[NodeId]struct{}) for _, location := range vl.vid2location { for _, dn := range location.list { if dcId := dn.GetDataCenterId(); dcId != "" { dataCenters[NodeId(dcId)] = struct{}{} if limit > 0 && len(dataCenters) >= limit { return dataCenters } } } } return dataCenters } // RackGrowPlan is one volume grow action produced by the periodic rack-aware // growth scan. An empty Rack means the grow is DC-wide. type RackGrowPlan struct { DataCenter string Rack string WritableVolumeCount uint32 } // PlanRackAwareGrowth returns the grow actions needed so every data center // this layout lives in keeps a non-crowded writable volume. stepCount is the // default per-event increment. // // For rack-spanning replication (DiffRackCount > 0) a single logical volume // already covers the racks the placement requires, so ShouldGrowVolumesByDcAndRack // returns the same result for every rack in a DC. Planning one grow per rack // would create racks×count too many volumes; plan one DC-wide grow instead. // The default increment is capped at the configured copy_N so lowering // master.volume_growth.copy_N reduces periodic growth. func (vl *VolumeLayout) PlanRackAwareGrowth(dcs map[NodeId][]NodeId, lastGrowCount, stepCount uint32) (plans []RackGrowPlan) { writables := vl.CloneWritableVolumes() if c := VolumeGrowthCountForCopies(vl.rp.GetCopyCount()); c < stepCount { stepCount = c } growOncePerDc := vl.rp.DiffRackCount > 0 // Only maintain data centers already hosting this layout's volumes: a // collection pinned to one DC (fs.configure -dataCenter) must not sprout // volumes in every other DC, and a DC that does need this layout gets its // volumes from the DC-constrained assign that first asks for them. hostingDCs := vl.listVolumeDataCenters(len(dcs)) // Spread lastGrowCount evenly across all grow targets. Summing every rack // up front keeps the divisor global, so DCs with different rack counts do // not each over-grow from a per-DC divisor. var rackPairs, dcCount uint32 for dcId, racks := range dcs { if _, hosting := hostingDCs[dcId]; !hosting { continue } dcCount++ rackPairs += uint32(len(racks)) } for dcId, racks := range dcs { if _, hosting := hostingDCs[dcId]; !hosting { continue } if growOncePerDc { if !vl.ShouldGrowVolumesByDcAndRack(&writables, dcId, "") { continue } count := stepCount if lastGrowCount > 0 { count = ceilDiv(lastGrowCount, dcCount) } plans = append(plans, RackGrowPlan{DataCenter: string(dcId), WritableVolumeCount: count}) continue } for _, rackId := range racks { if !vl.ShouldGrowVolumesByDcAndRack(&writables, dcId, rackId) { continue } count := stepCount if lastGrowCount > 0 { count = ceilDiv(lastGrowCount, rackPairs) } plans = append(plans, RackGrowPlan{DataCenter: string(dcId), Rack: string(rackId), WritableVolumeCount: count}) } } return plans } func ceilDiv(a, b uint32) uint32 { if b == 0 { return 0 } return (a + b - 1) / b } func (vl *VolumeLayout) GetWritableVolumeCount() (active, crowded int) { vl.accessLock.RLock() defer vl.accessLock.RUnlock() // The crowded map retains volumes that later became unwritable (full, // read-only), so their state survives transient writability flips. Count // only the writable ones: growth decisions compare crowded against // writables, and a raw len(vl.crowded) can exceed len(vl.writables) // permanently, demanding growth forever. for _, vid := range vl.writables { if _, ok := vl.crowded[vid]; ok { crowded++ } } return len(vl.writables), crowded } func (vl *VolumeLayout) CloneWritableVolumes() (writables []needle.VolumeId) { vl.accessLock.RLock() writables = make([]needle.VolumeId, len(vl.writables)) copy(writables, vl.writables) vl.accessLock.RUnlock() return writables } // CountUnderReplicatedVolumes returns the number of volumes in this layout // that do not have enough replicas according to their replica placement // configuration. Safe for concurrent access (RLock). func (vl *VolumeLayout) CountUnderReplicatedVolumes() int { vl.accessLock.RLock() defer vl.accessLock.RUnlock() count := 0 for vid := range vl.vid2location { if !vl.enoughCopies(vid) { count++ } } return count } func (vl *VolumeLayout) removeFromWritable(vid needle.VolumeId) bool { toDeleteIndex := -1 for k, id := range vl.writables { if id == vid { toDeleteIndex = k break } } if toDeleteIndex >= 0 { glog.V(0).Infoln("Volume", vid, "becomes unwritable") vl.writables = append(vl.writables[0:toDeleteIndex], vl.writables[toDeleteIndex+1:]...) return true } return false } func (vl *VolumeLayout) setVolumeWritable(vid needle.VolumeId) bool { for _, v := range vl.writables { if v == vid { return false } } glog.V(1).Infoln("Volume", vid, "becomes writable") vl.writables = append(vl.writables, vid) return true } func (vl *VolumeLayout) SetVolumeReadOnly(dn *DataNode, vid needle.VolumeId) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock() if location, ok := vl.vid2location[vid]; ok { location.SetReadOnly(dn, true) return vl.removeFromWritable(vid) } return true } func (vl *VolumeLayout) SetVolumeWritable(dn *DataNode, vid needle.VolumeId) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock() if location, ok := vl.vid2location[vid]; ok { location.SetReadOnly(dn, false) } if vl.enoughCopies(vid) { return vl.setVolumeWritable(vid) } return false } func (vl *VolumeLayout) SetVolumeUnavailable(dn *DataNode, vid needle.VolumeId) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock() if location, ok := vl.vid2location[vid]; ok { if removed := location.Remove(dn); removed != nil { moveLookupOwnership(vid, removed, nil) wasWritable := false if location.Length() < vl.rp.GetCopyCount() { glog.V(0).Infoln("Volume", vid, "has", location.Length(), "replica, less than required", vl.rp.GetCopyCount()) wasWritable = vl.removeFromWritable(vid) } if location.Length() == 0 { // Drop the now-empty entry. Otherwise Lookup returns a non-nil // empty location list, which surfaces as "volume id not found" // even though the volume still appears in volume.list/admin UI. // Mirrors UnRegisterVolume. delete(vl.vid2location, vid) delete(vl.sizeTracking, vid) vl.removeFromCrowded(vid) } return wasWritable } } return false } func (vl *VolumeLayout) SetVolumeAvailable(dn *DataNode, vid needle.VolumeId, isReadOnly, isFullCapacity bool) (becameWritable bool) { restoreActiveVolumeCount := false vl.accessLock.Lock() defer func() { vl.accessLock.Unlock() if restoreActiveVolumeCount { vl.adjustActiveVolumeCount(vid, +1) } }() vInfo, err := dn.GetVolumesById(vid) if err != nil { return false } if vl.dropped { return false } // A disconnect during a long vacuum can drop the entry while the volume is // still on the node; re-create it instead of dereferencing a nil location, // so the commit also repairs the split. moveLookupOwnership(vid, vl.getOrCreateLocationList(vid).Set(dn), dn) if vInfo.ReadOnly || isReadOnly || isFullCapacity { return false } vl.initSizeTracking(vid, vInfo.Size, vInfo.CompactRevision) if vl.enoughCopies(vid) { becameWritable = vl.setVolumeWritable(vid) if becameWritable { if st := vl.sizeTracking[vid]; st != nil && !st.fullSince.IsZero() { // fullSince marks a prior capacity-full removal that already // decremented activeVolumeCount. Re-adding must pair it once. st.fullSince = time.Time{} restoreActiveVolumeCount = true } } } return becameWritable } func (vl *VolumeLayout) enoughCopies(vid needle.VolumeId) bool { locations := vl.vid2location[vid].Length() desired := vl.rp.GetCopyCount() return locations == desired || (vl.replicationAsMin && locations > desired) } func (vl *VolumeLayout) SetVolumeCapacityFull(vid needle.VolumeId) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock() wasWritable := vl.removeFromWritable(vid) if wasWritable { // Stamp fullSince so UpdateVolumeSize's recovery branch can re-add the // volume once it shrinks; RecordAssign does this for the write path. // Only on actual removal, to stay paired with the activeVolumeCount // decrement the caller does for the same bool. if st := vl.sizeTracking[vid]; st != nil && st.fullSince.IsZero() { st.fullSince = time.Now() } glog.V(0).Infof("Volume %d reaches full capacity.", vid) } return wasWritable } func (vl *VolumeLayout) removeFromCrowded(vid needle.VolumeId) { if _, ok := vl.crowded[vid]; ok { glog.V(0).Infoln("Volume", vid, "becomes uncrowded") delete(vl.crowded, vid) } } func (vl *VolumeLayout) setVolumeCrowded(vid needle.VolumeId) { if _, ok := vl.crowded[vid]; !ok { vl.crowded[vid] = struct{}{} glog.V(0).Infoln("Volume", vid, "becomes crowded") } } func (vl *VolumeLayout) SetVolumeCrowded(vid needle.VolumeId) { // since delete is guarded by accessLock.Lock(), // and is always called in sequential order, // RLock() should be safe enough vl.accessLock.RLock() defer vl.accessLock.RUnlock() vl.setVolumeCrowded(vid) } type VolumeLayoutInfo struct { Replication string `json:"replication"` TTL string `json:"ttl"` Writables []needle.VolumeId `json:"writables"` Collection string `json:"collection"` DiskType string `json:"diskType"` } func (vl *VolumeLayout) ToInfo() (info VolumeLayoutInfo) { info.Replication = vl.rp.String() info.TTL = vl.ttl.String() info.Writables = vl.writables info.DiskType = vl.diskType.ReadableString() //m["locations"] = vl.vid2location return } func (vlc *VolumeLayoutCollection) ToVolumeGrowRequest() *master_pb.VolumeGrowRequest { vgr := &master_pb.VolumeGrowRequest{ Collection: vlc.Collection, Replication: vlc.VolumeLayout.rp.String(), Ttl: vlc.VolumeLayout.ttl.String(), DiskType: vlc.VolumeLayout.diskType.String(), } // A layout living in exactly one data center is either pinned there or on // a single-DC cluster; either way unconstrained growth has no business // placing its volumes elsewhere. Cross-DC replication can never // legitimately live in one DC — there a single hosting DC means an // outage, not a pin. if vlc.VolumeLayout.rp.DiffDataCenterCount > 0 { return vgr } if dcs := vlc.VolumeLayout.listVolumeDataCenters(2); len(dcs) == 1 { for dc := range dcs { vgr.DataCenter = string(dc) } } return vgr } func (vl *VolumeLayout) Stats() *VolumeLayoutStats { vl.accessLock.RLock() defer vl.accessLock.RUnlock() ret := &VolumeLayoutStats{} freshThreshold := time.Now().Unix() - 60 for vid, vll := range vl.vid2location { size, fileCount := vll.Stats(vid, freshThreshold) ret.FileCount += uint64(fileCount) ret.UsedSize += size * uint64(vll.Length()) ret.LogicalUsedSize += size if vll.AnyReadOnly() { ret.TotalSize += size * uint64(vll.Length()) } else { ret.TotalSize += vl.volumeSizeLimit * uint64(vll.Length()) } } return ret }