topology: keep per-node volume state with the location it describes (#10654)

* topology: keep per-node volume state with the location it describes

The read-only and oversized indexes were maps from volume id to a list of the
nodes reporting that state -- the same key space the lookup index already
holds, kept a second and third time. Nothing ever asked which nodes; both are
only ever asked whether any node does.

So the state rides on the location list as a bit per entry, and the two indexes
go. Removing a location shifts the bits with it, and a node replacing another
at the same address inherits its slot, since that is what happens to the
location too.

800k volumes, 90% read-only: readonly index 40.0MB -> 0, lookup index
42.5MB -> 48.6MB for the bits, 33.9MB net.

* topology: rebuild the location flags when stale entries are dropped

Refresh rebuilds the location list, so leaving the flags alone left bits
describing whoever moved into the dropped entries' place.

* topology: assert the refreshed flag survived, not just that it moved

Clearing the mask rather than rebuilding it would have passed: the check that
the flag lands on the right location is done by clearing it, which an already
empty mask satisfies.
This commit is contained in:
Chris Lu
2026-08-08 20:56:08 -07:00
committed by GitHub
parent 2dc59c9b51
commit 67b0cc0706
4 changed files with 220 additions and 101 deletions
+1 -1
View File
@@ -344,7 +344,7 @@ func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeL
// it so a benignly read-only (full/oversized) volume can be reclaimed.
func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayout *VolumeLayout, c *Collection, garbageThreshold float64, locationList *VolumeLocationList, vid needle.VolumeId, preallocate int64, skipReadOnly bool) {
volumeLayout.accessLock.RLock()
isReadOnly := volumeLayout.readonlyVolumes.IsTrue(vid)
isReadOnly := volumeLayout.vid2location[vid].AnyReadOnly()
isEnoughCopies := volumeLayout.enoughCopies(vid)
volumeLayout.accessLock.RUnlock()
+14 -91
View File
@@ -35,73 +35,6 @@ const (
NoWritableVolumes = "No writable volumes"
)
type stateIndicator func(copyState) bool
func ExistCopies() stateIndicator {
return func(state copyState) bool { return state != noCopies }
}
type volumesBinaryState struct {
rp *super_block.ReplicaPlacement
name volumeState // the name for volume state (eg. "Readonly", "Oversized")
indicator stateIndicator // indicate whether the volumes should be marked as `name`
copyMap map[needle.VolumeId]*VolumeLocationList
}
func NewVolumesBinaryState(name volumeState, rp *super_block.ReplicaPlacement, indicator stateIndicator) *volumesBinaryState {
return &volumesBinaryState{
rp: rp,
name: name,
indicator: indicator,
copyMap: make(map[needle.VolumeId]*VolumeLocationList),
}
}
func (v *volumesBinaryState) Dump() (res []uint32) {
for vid, list := range v.copyMap {
if v.indicator(v.copyState(list)) {
res = append(res, uint32(vid))
}
}
return
}
func (v *volumesBinaryState) IsTrue(vid needle.VolumeId) bool {
list, _ := v.copyMap[vid]
return v.indicator(v.copyState(list))
}
func (v *volumesBinaryState) Add(vid needle.VolumeId, dn *DataNode) {
list, _ := v.copyMap[vid]
if list != nil {
list.Set(dn)
return
}
list = NewVolumeLocationList()
list.Set(dn)
v.copyMap[vid] = list
}
func (v *volumesBinaryState) Remove(vid needle.VolumeId, dn *DataNode) {
list, _ := v.copyMap[vid]
if list != nil {
list.Remove(dn)
if list.Length() == 0 {
delete(v.copyMap, vid)
}
}
}
func (v *volumesBinaryState) copyState(list *VolumeLocationList) copyState {
if list == nil {
return noCopies
}
if list.Length() < v.rp.GetCopyCount() {
return insufficientCopies
}
return enoughCopies
}
// volumeSizeTracking holds per-volume size accounting for weighted assignment.
type volumeSizeTracking struct {
effectiveSize uint64 // reported + pending assigned bytes
@@ -128,8 +61,6 @@ type VolumeLayout struct {
vid2location map[needle.VolumeId]*VolumeLocationList
writables []needle.VolumeId // transient array of writable volume id
crowded map[needle.VolumeId]struct{}
readonlyVolumes *volumesBinaryState // readonly volumes
oversizedVolumes *volumesBinaryState // oversized volumes
vacuumedVolumes map[needle.VolumeId]time.Time
volumeSizeLimit uint64
replicationAsMin bool
@@ -154,8 +85,6 @@ func NewVolumeLayout(rp *super_block.ReplicaPlacement, ttl *needle.TTL, diskType
vid2location: make(map[needle.VolumeId]*VolumeLocationList),
writables: *new([]needle.VolumeId),
crowded: make(map[needle.VolumeId]struct{}),
readonlyVolumes: NewVolumesBinaryState(readOnlyState, rp, ExistCopies()),
oversizedVolumes: NewVolumesBinaryState(oversizedState, rp, ExistCopies()),
vacuumedVolumes: make(map[needle.VolumeId]time.Time),
volumeSizeLimit: volumeSizeLimit,
replicationAsMin: replicationAsMin,
@@ -201,20 +130,20 @@ func (vl *VolumeLayout) RegisterVolume(v *storage.VolumeInfo, dn *DataNode) {
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())
for _, dn := range vl.vid2location[v.Id].list {
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)
vl.readonlyVolumes.Add(v.Id, dn)
location.SetReadOnly(dn, true)
return
} else {
vl.readonlyVolumes.Remove(v.Id, dn)
}
location.SetReadOnly(dn, false)
} else {
glog.V(1).Infof("vid %d removed from writable", v.Id)
vl.removeFromWritable(v.Id)
vl.readonlyVolumes.Remove(v.Id, dn)
location.SetReadOnly(dn, false)
return
}
}
@@ -222,10 +151,8 @@ func (vl *VolumeLayout) RegisterVolume(v *storage.VolumeInfo, dn *DataNode) {
}
func (vl *VolumeLayout) rememberOversizedVolume(v *storage.VolumeInfo, dn *DataNode) {
if vl.isOversized(v) {
vl.oversizedVolumes.Add(v.Id, dn)
} else {
vl.oversizedVolumes.Remove(v.Id, dn)
if location, ok := vl.vid2location[v.Id]; ok {
location.SetOversized(dn, vl.isOversized(v))
}
}
@@ -305,7 +232,7 @@ func (vl *VolumeLayout) UpdateVolumeSize(vid needle.VolumeId, reportedSize uint6
if reportedSize >= vl.volumeSizeLimit {
return false // actual on-disk size still over limit; stay out
}
if vl.oversizedVolumes.IsTrue(vid) {
if vl.vid2location[vid].AnyOversized() {
return false
}
if !vl.enoughCopies(vid) || !vl.isAllWritable(vid) {
@@ -333,8 +260,6 @@ func (vl *VolumeLayout) UnRegisterVolume(v *storage.VolumeInfo, dn *DataNode) {
if removed := location.Remove(dn); removed != nil {
moveLookupOwnership(v.Id, removed, nil)
vl.readonlyVolumes.Remove(v.Id, dn)
vl.oversizedVolumes.Remove(v.Id, dn)
vl.ensureCorrectWritables(v.Id)
if location.Length() == 0 {
@@ -356,7 +281,7 @@ func (vl *VolumeLayout) EnsureCorrectWritables(v *storage.VolumeInfo) {
func (vl *VolumeLayout) ensureCorrectWritables(vid needle.VolumeId) {
isEnoughCopies := vl.enoughCopies(vid)
isAllWritable := vl.isAllWritable(vid)
isOversizedVolume := vl.oversizedVolumes.IsTrue(vid)
isOversizedVolume := vl.vid2location[vid].AnyOversized()
if isEnoughCopies && isAllWritable && !isOversizedVolume {
vl.setVolumeWritable(vid)
return
@@ -869,8 +794,8 @@ func (vl *VolumeLayout) SetVolumeReadOnly(dn *DataNode, vid needle.VolumeId) boo
vl.accessLock.Lock()
defer vl.accessLock.Unlock()
if _, ok := vl.vid2location[vid]; ok {
vl.readonlyVolumes.Add(vid, dn)
if location, ok := vl.vid2location[vid]; ok {
location.SetReadOnly(dn, true)
return vl.removeFromWritable(vid)
}
return true
@@ -880,8 +805,8 @@ func (vl *VolumeLayout) SetVolumeWritable(dn *DataNode, vid needle.VolumeId) boo
vl.accessLock.Lock()
defer vl.accessLock.Unlock()
if _, ok := vl.vid2location[vid]; ok {
vl.readonlyVolumes.Remove(vid, dn)
if location, ok := vl.vid2location[vid]; ok {
location.SetReadOnly(dn, false)
}
if vl.enoughCopies(vid) {
@@ -897,8 +822,6 @@ func (vl *VolumeLayout) SetVolumeUnavailable(dn *DataNode, vid needle.VolumeId)
if location, ok := vl.vid2location[vid]; ok {
if removed := location.Remove(dn); removed != nil {
moveLookupOwnership(vid, removed, nil)
vl.readonlyVolumes.Remove(vid, dn)
vl.oversizedVolumes.Remove(vid, dn)
wasWritable := false
if location.Length() < vl.rp.GetCopyCount() {
glog.V(0).Infoln("Volume", vid, "has", location.Length(), "replica, less than required", vl.rp.GetCopyCount())
@@ -1044,7 +967,7 @@ func (vl *VolumeLayout) Stats() *VolumeLayoutStats {
ret.FileCount += uint64(fileCount)
ret.UsedSize += size * uint64(vll.Length())
ret.LogicalUsedSize += size
if vl.readonlyVolumes.IsTrue(vid) {
if vll.AnyReadOnly() {
ret.TotalSize += size * uint64(vll.Length())
} else {
ret.TotalSize += vl.volumeSizeLimit * uint64(vll.Length())
+120
View File
@@ -0,0 +1,120 @@
package topology
import (
"testing"
)
func flagNode(ip string) *DataNode {
dn := NewDataNode(ip)
dn.Ip, dn.Port = ip, 8080
return dn
}
func TestLocationFlagsFollowTheirNode(t *testing.T) {
a, b, c := flagNode("10.0.0.1"), flagNode("10.0.0.2"), flagNode("10.0.0.3")
list := NewVolumeLocationList()
for _, dn := range []*DataNode{a, b, c} {
list.Set(dn)
}
list.SetReadOnly(b, true)
list.SetOversized(c, true)
if !list.AnyReadOnly() || !list.AnyOversized() {
t.Fatal("flags did not register")
}
// Removing an earlier node shifts the rest down; the flags have to move with
// them or they end up describing the wrong server.
list.Remove(a)
list.SetReadOnly(b, false)
if list.AnyReadOnly() {
t.Error("clearing the flag on its node left it set, so it had shifted onto another")
}
if !list.AnyOversized() {
t.Error("removing an unrelated node dropped another node's flag")
}
list.SetOversized(c, false)
if list.AnyOversized() {
t.Error("clearing the flag on its node left it set")
}
}
func TestLocationFlagsClearedWithTheirNode(t *testing.T) {
a, b := flagNode("10.0.0.1"), flagNode("10.0.0.2")
list := NewVolumeLocationList()
list.Set(a)
list.Set(b)
list.SetReadOnly(a, true)
list.Remove(a)
if list.AnyReadOnly() {
t.Error("a departed node left its read-only marking behind")
}
}
// A node that shares an address replaces the entry, so it inherits the slot.
func TestLocationFlagsSurviveAReplacedNode(t *testing.T) {
a, replacement := flagNode("10.0.0.1"), flagNode("10.0.0.1")
list := NewVolumeLocationList()
list.Set(a)
list.SetReadOnly(a, true)
list.Set(replacement)
if !list.AnyReadOnly() {
t.Error("replacing the node at an address dropped what was known about the volume there")
}
list.SetReadOnly(replacement, false)
if list.AnyReadOnly() {
t.Error("the replacement could not clear the flag it inherited")
}
}
func TestLocationFlagsIgnoreUntrackableReplicaCounts(t *testing.T) {
list := NewVolumeLocationList()
nodes := make([]*DataNode, 0, maxTrackedLocations+2)
for i := 0; i < maxTrackedLocations+2; i++ {
dn := NewDataNode("n")
dn.Ip, dn.Port = "10.0.0.1", 9000+i
nodes = append(nodes, dn)
list.Set(dn)
}
// Past the tracked width the flag is not recorded rather than landing on
// some other node's slot.
list.SetReadOnly(nodes[maxTrackedLocations+1], true)
if list.AnyReadOnly() {
t.Error("a flag beyond the tracked width was recorded against another node")
}
list.SetReadOnly(nodes[0], true)
if !list.AnyReadOnly() {
t.Error("a flag within the tracked width was lost")
}
}
// Refresh drops stale locations, so it has to rebuild the flags with them.
func TestLocationFlagsRebuiltByRefresh(t *testing.T) {
stale, fresh, alsoFresh := flagNode("10.0.0.1"), flagNode("10.0.0.2"), flagNode("10.0.0.3")
stale.LastSeen, fresh.LastSeen, alsoFresh.LastSeen = 100, 500, 500
list := NewVolumeLocationList()
for _, dn := range []*DataNode{stale, fresh, alsoFresh} {
list.Set(dn)
}
list.SetReadOnly(alsoFresh, true)
list.Refresh(400)
if list.Length() != 2 {
t.Fatalf("expected the stale location to be dropped, got %d", list.Length())
}
// Both halves are needed: that the flag survived at all, and that it
// survived against the right location. Rebuilding the mask as empty would
// pass the second on its own.
if !list.AnyReadOnly() {
t.Fatal("the flag was dropped rather than carried across the rebuild")
}
list.SetReadOnly(alsoFresh, false)
if list.AnyReadOnly() {
t.Error("the flag did not move with its location, so it now describes another server")
}
}
+85 -9
View File
@@ -6,8 +6,18 @@ import (
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
)
// maxTrackedLocations is how many replicas of one volume can carry per-node
// state. Beyond that the flags are simply not recorded, which loses a read-only
// or oversized marking rather than misreporting one.
const maxTrackedLocations = 32
type VolumeLocationList struct {
list []*DataNode
// readOnly and oversized mirror list by index, so the per-node state that
// used to need an index of its own rides along with the location it
// describes. Both are only ever asked whether any node reports them.
readOnly uint32
oversized uint32
}
func NewVolumeLocationList() *VolumeLocationList {
@@ -22,10 +32,67 @@ func (dnll *VolumeLocationList) Copy() *VolumeLocationList {
list := make([]*DataNode, len(dnll.list))
copy(list, dnll.list)
return &VolumeLocationList{
list: list,
list: list,
readOnly: dnll.readOnly,
oversized: dnll.oversized,
}
}
func (dnll *VolumeLocationList) indexOf(loc *DataNode) int {
for i, dnl := range dnll.list {
if loc.Ip == dnl.Ip && loc.Port == dnl.Port {
return i
}
}
return -1
}
func setFlag(flags *uint32, index int, on bool) {
if index < 0 || index >= maxTrackedLocations {
return
}
if on {
*flags |= 1 << uint(index)
} else {
*flags &^= 1 << uint(index)
}
}
func hasFlag(flags uint32, index int) bool {
return index >= 0 && index < maxTrackedLocations && flags&(1<<uint(index)) != 0
}
// removeFlag drops the bit at index and closes the gap, keeping the flags
// aligned with the locations after one is removed.
func removeFlag(flags *uint32, index int) {
if index < 0 || index >= maxTrackedLocations {
return
}
low := *flags & (1<<uint(index) - 1)
high := *flags >> uint(index+1) << uint(index)
*flags = low | high
}
// SetReadOnly records whether loc reports the volume read-only.
func (dnll *VolumeLocationList) SetReadOnly(loc *DataNode, readOnly bool) {
setFlag(&dnll.readOnly, dnll.indexOf(loc), readOnly)
}
// AnyReadOnly reports whether any location has the volume read-only.
func (dnll *VolumeLocationList) AnyReadOnly() bool {
return dnll != nil && dnll.readOnly != 0
}
// SetOversized records whether loc reports the volume past the size limit.
func (dnll *VolumeLocationList) SetOversized(loc *DataNode, oversized bool) {
setFlag(&dnll.oversized, dnll.indexOf(loc), oversized)
}
// AnyOversized reports whether any location has the volume past the size limit.
func (dnll *VolumeLocationList) AnyOversized() bool {
return dnll != nil && dnll.oversized != 0
}
func (dnll *VolumeLocationList) Head() *DataNode {
//mark first node as master volume
if dnll.Length() == 0 {
@@ -71,6 +138,8 @@ func (dnll *VolumeLocationList) Remove(loc *DataNode) (removed *DataNode) {
for i, dnl := range dnll.list {
if loc.Ip == dnl.Ip && loc.Port == dnl.Port {
dnll.list = append(dnll.list[:i], dnll.list[i+1:]...)
removeFlag(&dnll.readOnly, i)
removeFlag(&dnll.oversized, i)
return dnl
}
}
@@ -85,15 +154,22 @@ func (dnll *VolumeLocationList) Refresh(freshThreshHold int64) {
break
}
}
if changed {
var l []*DataNode
for _, dnl := range dnll.list {
if dnl.LastSeen >= freshThreshHold {
l = append(l, dnl)
}
}
dnll.list = l
if !changed {
return
}
// The flags index the locations, so dropping some means rebuilding both
// rather than leaving bits describing whoever moved into their place.
var l []*DataNode
var readOnly, oversized uint32
for i, dnl := range dnll.list {
if dnl.LastSeen < freshThreshHold {
continue
}
setFlag(&readOnly, len(l), hasFlag(dnll.readOnly, i))
setFlag(&oversized, len(l), hasFlag(dnll.oversized, i))
l = append(l, dnl)
}
dnll.list, dnll.readOnly, dnll.oversized = l, readOnly, oversized
}
// Stats returns logic size and count. Both subtractions are clamped: the