diff --git a/weed/topology/topology.go b/weed/topology/topology.go index 1ac1edcf3..847379249 100644 --- a/weed/topology/topology.go +++ b/weed/topology/topology.go @@ -651,6 +651,12 @@ func (t *Topology) SyncDataNodeRegistration(volumes []*master_pb.VolumeInformati for _, v := range deletedVolumes { t.UnRegisterVolumeLayout(v, dn) } + for _, v := range changedVolumes { + if v.ReplicaPlacement == nil { + continue + } + t.GetVolumeLayout(v.Collection, v.ReplicaPlacement, v.Ttl, types.ToDiskType(v.DiskType)).SetReplicaReadOnlyFlag(dn, v.Id, v.ReadOnly) + } // Update effective sizes for all reported volumes (decay pending estimates). // If decay brings a volume eagerly removed by RecordAssign back under the // writable threshold, restore the matching activeVolumeCount. @@ -780,6 +786,9 @@ func (t *Topology) ApplyVolumeChanges(changed []*master_pb.VolumeInformationMess if isNew || becameServable || tierTransition || isChanged { newVolumes = append(newVolumes, vi) } + if isChanged { + vl.SetReplicaReadOnlyFlag(dn, vi.Id, vi.ReadOnly) + } vl.UpdateOversizedState(&vi, dn) if vl.UpdateVolumeSize(vi.Id, vi.Size, vi.CompactRevision, true) { vl.AdjustActiveVolumeCountAfterRecovery(vi.Id) diff --git a/weed/topology/topology_test.go b/weed/topology/topology_test.go index bb4054568..3a89ce14f 100644 --- a/weed/topology/topology_test.go +++ b/weed/topology/topology_test.go @@ -864,3 +864,48 @@ func TestSetVolumeAvailableRepairsMissingVolume(t *testing.T) { t.Fatalf("after SetVolumeAvailable: size tracking for %d not seeded", vid) } } + +// A read-only change that arrives in the regular heartbeat must reach the +// layout's per-replica flag, which the vacuum sweep consults. Registration +// and volume.mark already set it; the heartbeat path did not. +func TestFullHeartbeatUpdatesLayoutReadOnlyFlag(t *testing.T) { + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + dn := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1"). + GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", map[string]uint32{"": 25}) + report := func(readOnly bool) []*master_pb.VolumeInformationMessage { + return []*master_pb.VolumeInformationMessage{{ + Id: 1, Collection: "c", Size: 1 << 20, Version: uint32(needle.GetCurrentVersion()), ReadOnly: readOnly, + }} + } + rp, _ := super_block.NewReplicaPlacementFromString("000") + flag := func() bool { + vl := topo.GetVolumeLayout("c", rp, needle.EMPTY_TTL, types.HardDriveType) + vl.accessLock.RLock() + defer vl.accessLock.RUnlock() + return vl.vid2location[needle.VolumeId(1)].AnyReadOnly() + } + + // the full list, which a server sends first and whenever the digests disagree + topo.SyncDataNodeRegistration(report(false), dn) + if flag() { + t.Fatal("volume registered writable is flagged read-only") + } + topo.SyncDataNodeRegistration(report(true), dn) + if !flag() { + t.Fatal("full heartbeat turned the volume read-only but the layout flag did not follow") + } + topo.SyncDataNodeRegistration(report(false), dn) + if flag() { + t.Fatal("full heartbeat turned the volume writable again but the layout flag stayed read-only") + } + + // the changed-volumes delta, which is what a running server normally sends + topo.ApplyVolumeChanges(report(true), dn) + if !flag() { + t.Fatal("delta heartbeat turned the volume read-only but the layout flag did not follow") + } + topo.ApplyVolumeChanges(report(false), dn) + if flag() { + t.Fatal("delta heartbeat turned the volume writable again but the layout flag stayed read-only") + } +} diff --git a/weed/topology/volume_layout.go b/weed/topology/volume_layout.go index b8da55d5f..31847adb5 100644 --- a/weed/topology/volume_layout.go +++ b/weed/topology/volume_layout.go @@ -988,6 +988,21 @@ func (vl *VolumeLayout) SetVolumeWritable(dn *DataNode, vid needle.VolumeId) boo return false } +// SetReplicaReadOnlyFlag records what dn's latest heartbeat said about its +// replica of vid, and nothing more: the writable list is left to +// EnsureCorrectWritables and its capacity guards. Until now the flag only moved +// on registration and volume.mark, so a replica that went read-only, or came +// back, while the server ran kept its stale flag until the next restart. The +// vacuum sweep reads this flag when it decides whether to skip a volume. +func (vl *VolumeLayout) SetReplicaReadOnlyFlag(dn *DataNode, vid needle.VolumeId, readOnly bool) { + vl.accessLock.Lock() + defer vl.accessLock.Unlock() + + if location, ok := vl.vid2location[vid]; ok { + location.SetReadOnly(dn, readOnly) + } +} + func (vl *VolumeLayout) SetVolumeUnavailable(dn *DataNode, vid needle.VolumeId) bool { vl.accessLock.Lock() defer vl.accessLock.Unlock()