diff --git a/weed/command/scaffold/master.toml b/weed/command/scaffold/master.toml index 38443410b..34a7168c6 100644 --- a/weed/command/scaffold/master.toml +++ b/weed/command/scaffold/master.toml @@ -41,6 +41,7 @@ copy_2 = 6 # create 2 x 6 = 12 actual volumes copy_3 = 3 # create 3 x 3 = 9 actual volumes copy_other = 1 # create n x 1 = n actual volumes threshold = 0.9 # create threshold +reservation_timeout = "5m"# capacity reservation timeout before unreleased reservations expire disable = false # disables volume growth if true # configuration flags for replication diff --git a/weed/server/master_server.go b/weed/server/master_server.go index dbf746997..22f33176e 100644 --- a/weed/server/master_server.go +++ b/weed/server/master_server.go @@ -119,6 +119,7 @@ func NewMasterServer(r *mux.Router, option *MasterOption, peers map[string]pb.Se v.SetDefault("master.volume_growth.copy_3", topology.VolumeGrowStrategy.Copy3Count) v.SetDefault("master.volume_growth.copy_other", topology.VolumeGrowStrategy.CopyOtherCount) v.SetDefault("master.volume_growth.threshold", topology.VolumeGrowStrategy.Threshold) + v.SetDefault("master.volume_growth.reservation_timeout", "5m") v.SetDefault("master.volume_growth.disable", false) option.VolumeGrowthDisabled = v.GetBool("master.volume_growth.disable") @@ -127,6 +128,7 @@ func NewMasterServer(r *mux.Router, option *MasterOption, peers map[string]pb.Se topology.VolumeGrowStrategy.Copy3Count = v.GetUint32("master.volume_growth.copy_3") topology.VolumeGrowStrategy.CopyOtherCount = v.GetUint32("master.volume_growth.copy_other") topology.VolumeGrowStrategy.Threshold = v.GetFloat64("master.volume_growth.threshold") + topology.VolumeGrowStrategy.ReservationTimeout = parseReservationTimeout(v) whiteList := util.StringSplit(v.GetString("guard.white_list"), ",") var preallocateSize int64 @@ -663,3 +665,16 @@ func (ms *MasterServer) Reload() { v.GetInt("jwt.signing.read.expires_after_seconds"), ) } + +func parseReservationTimeout(v *util.ViperProxy) time.Duration { + if str := strings.TrimSpace(v.GetString("master.volume_growth.reservation_timeout")); str != "" { + if d, err := time.ParseDuration(str); err == nil && d > 0 { + return d + } + } + // A bare number in the config means seconds; GetDuration would read it as nanoseconds. + if sec := v.GetInt("master.volume_growth.reservation_timeout"); sec > 0 { + return time.Duration(sec) * time.Second + } + return topology.DefaultReservationTimeout +} diff --git a/weed/server/master_server_reservation_timeout_test.go b/weed/server/master_server_reservation_timeout_test.go new file mode 100644 index 000000000..ca9aecc82 --- /dev/null +++ b/weed/server/master_server_reservation_timeout_test.go @@ -0,0 +1,85 @@ +package weed_server + +import ( + "testing" + "time" + + "github.com/spf13/viper" + + "github.com/seaweedfs/seaweedfs/weed/util" +) + +func TestParseReservationTimeout(t *testing.T) { + tests := []struct { + name string + setup func(v *util.ViperProxy) + expected time.Duration + }{ + { + name: "unset defaults to 5 minutes", + setup: func(v *util.ViperProxy) {}, + expected: 5 * time.Minute, + }, + { + name: "duration string 10m", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", "10m") + }, + expected: 10 * time.Minute, + }, + { + name: "duration string 30s", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", "30s") + }, + expected: 30 * time.Second, + }, + { + name: "integer seconds 300", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", 300) + }, + expected: 300 * time.Second, + }, + { + name: "large integer stays seconds not nanoseconds", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", 1000000000) + }, + expected: 1000000000 * time.Second, + }, + { + name: "non-positive zero defaults to 5 minutes", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", 0) + }, + expected: 5 * time.Minute, + }, + { + name: "negative seconds defaults to 5 minutes", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", -10) + }, + expected: 5 * time.Minute, + }, + { + name: "invalid string defaults to 5 minutes", + setup: func(v *util.ViperProxy) { + v.Set("master.volume_growth.reservation_timeout", "invalid") + }, + expected: 5 * time.Minute, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + rawViper := viper.New() + vp := util.NewViperProxy(rawViper) + tt.setup(vp) + actual := parseReservationTimeout(vp) + if actual != tt.expected { + t.Errorf("expected %v, got %v", tt.expected, actual) + } + }) + } +} diff --git a/weed/topology/node.go b/weed/topology/node.go index a603d7d87..b8b51eb48 100644 --- a/weed/topology/node.go +++ b/weed/topology/node.go @@ -349,17 +349,18 @@ func (n *NodeImpl) CapacityForAnyDisk() (total int64) { // AvailableSpaceForReservation returns available space considering existing reservations func (n *NodeImpl) AvailableSpaceForReservation(option *VolumeGrowOption) int64 { + // Expire here as well: a node whose reservations fill it is filtered out + // before TryReserveCapacity could clean them, stranding the capacity. + n.capacityReservations.cleanExpiredReservations(VolumeGrowStrategy.GetReservationTimeout()) baseAvailable := n.AvailableSpaceFor(option) reservedCount := n.capacityReservations.getReservedCount(option.DiskType) return baseAvailable - reservedCount } -// TryReserveCapacity attempts to atomically reserve capacity for volume creation +// TryReserveCapacity attempts to atomically reserve capacity for volume creation using the configured timeout func (n *NodeImpl) TryReserveCapacity(diskType types.DiskType, count int64) (reservationId string, success bool) { - const reservationTimeout = 5 * time.Minute // TODO: make this configurable - // Clean up any expired reservations first - n.capacityReservations.cleanExpiredReservations(reservationTimeout) + n.capacityReservations.cleanExpiredReservations(VolumeGrowStrategy.GetReservationTimeout()) // Atomically check and reserve space option := &VolumeGrowOption{DiskType: diskType} diff --git a/weed/topology/volume_growth.go b/weed/topology/volume_growth.go index c4aff30fe..e50b6a726 100644 --- a/weed/topology/volume_growth.go +++ b/weed/topology/volume_growth.go @@ -39,21 +39,32 @@ func (vg *VolumeGrowRequest) Equals(req *VolumeGrowRequest) bool { return reflect.DeepEqual(vg.Option, req.Option) && vg.Count == req.Count && vg.Force == req.Force } +const DefaultReservationTimeout = 5 * time.Minute + type volumeGrowthStrategy struct { - Copy1Count uint32 - Copy2Count uint32 - Copy3Count uint32 - CopyOtherCount uint32 - Threshold float64 + Copy1Count uint32 + Copy2Count uint32 + Copy3Count uint32 + CopyOtherCount uint32 + Threshold float64 + ReservationTimeout time.Duration +} + +func (s *volumeGrowthStrategy) GetReservationTimeout() time.Duration { + if s != nil && s.ReservationTimeout > 0 { + return s.ReservationTimeout + } + return DefaultReservationTimeout } var ( VolumeGrowStrategy = volumeGrowthStrategy{ - Copy1Count: 7, - Copy2Count: 6, - Copy3Count: 3, - CopyOtherCount: 1, - Threshold: 0.9, + Copy1Count: 7, + Copy2Count: 6, + Copy3Count: 3, + CopyOtherCount: 1, + Threshold: 0.9, + ReservationTimeout: DefaultReservationTimeout, } ) diff --git a/weed/topology/volume_growth_reservation_test.go b/weed/topology/volume_growth_reservation_test.go index f3711b2fa..08d7952c3 100644 --- a/weed/topology/volume_growth_reservation_test.go +++ b/weed/topology/volume_growth_reservation_test.go @@ -293,3 +293,111 @@ func TestVolumeGrowth_ReservationTimeout(t *testing.T) { t.Errorf("Expected 2 available slots after cleanup and new reservation, got %d", available) } } + +func TestVolumeGrowth_ConfigurableReservationTimeout(t *testing.T) { + origTimeout := VolumeGrowStrategy.ReservationTimeout + defer func() { + VolumeGrowStrategy.ReservationTimeout = origTimeout + }() + + dn := NewDataNode("server1") + diskType := types.HardDriveType + + // Set up capacity of 5 + diskUsage := dn.diskUsages.getOrCreateDisk(diskType) + diskUsage.maxVolumeCount = 5 + + // 1. Verify default timeout (5 minutes) + VolumeGrowStrategy.ReservationTimeout = 5 * time.Minute + resId1, ok := dn.TryReserveCapacity(diskType, 2) + if !ok { + t.Fatal("Expected reservation 1 to succeed") + } + + // Set reservation createdAt to 4 minutes ago (not expired under 5m timeout) + dn.capacityReservations.Lock() + if r, exists := dn.capacityReservations.reservations[resId1]; exists { + r.createdAt = time.Now().Add(-4 * time.Minute) + } + dn.capacityReservations.Unlock() + + // Available space should be 5 - 2 = 3. Trying to reserve 4 must fail. + _, ok = dn.TryReserveCapacity(diskType, 4) + if ok { + t.Error("Expected reservation of 4 to fail when 2 slots are still reserved") + } + + // Set reservation createdAt to 6 minutes ago (expired under 5m timeout) + dn.capacityReservations.Lock() + if r, exists := dn.capacityReservations.reservations[resId1]; exists { + r.createdAt = time.Now().Add(-6 * time.Minute) + } + dn.capacityReservations.Unlock() + + // Now reserving 4 should clean up the expired reservation and succeed + resId2, ok := dn.TryReserveCapacity(diskType, 4) + if !ok { + t.Fatal("Expected reservation of 4 to succeed after 6m expired reservation was cleaned up") + } + dn.ReleaseReservedCapacity(resId2) + + // 2. Verify custom timeout (1 minute) + VolumeGrowStrategy.ReservationTimeout = 1 * time.Minute + resId3, ok := dn.TryReserveCapacity(diskType, 2) + if !ok { + t.Fatal("Expected reservation 3 to succeed") + } + + // Set createdAt to 45 seconds ago (not expired under 1m timeout) + dn.capacityReservations.Lock() + if r, exists := dn.capacityReservations.reservations[resId3]; exists { + r.createdAt = time.Now().Add(-45 * time.Second) + } + dn.capacityReservations.Unlock() + + _, ok = dn.TryReserveCapacity(diskType, 4) + if ok { + t.Error("Expected reservation of 4 to fail when 2 slots are reserved 45s ago with 1m timeout") + } + + // Set createdAt to 75 seconds ago (expired under 1m timeout) + dn.capacityReservations.Lock() + if r, exists := dn.capacityReservations.reservations[resId3]; exists { + r.createdAt = time.Now().Add(-75 * time.Second) + } + dn.capacityReservations.Unlock() + + resId4, ok := dn.TryReserveCapacity(diskType, 4) + if !ok { + t.Fatal("Expected reservation of 4 to succeed after 75s reservation expired under 1m timeout") + } + dn.ReleaseReservedCapacity(resId4) + + // 3. Verify non-positive timeout fallback to 5 minutes + VolumeGrowStrategy.ReservationTimeout = 0 + if VolumeGrowStrategy.GetReservationTimeout() != 5*time.Minute { + t.Errorf("Expected 0 timeout to fall back to 5m, got %v", VolumeGrowStrategy.GetReservationTimeout()) + } + VolumeGrowStrategy.ReservationTimeout = -10 * time.Second + if VolumeGrowStrategy.GetReservationTimeout() != 5*time.Minute { + t.Errorf("Expected negative timeout to fall back to 5m, got %v", VolumeGrowStrategy.GetReservationTimeout()) + } + + // 4. Expired reservations must not strand capacity: the selection filter + // reads AvailableSpaceForReservation without calling TryReserveCapacity. + VolumeGrowStrategy.ReservationTimeout = 1 * time.Minute + resId5, ok := dn.TryReserveCapacity(diskType, 5) + if !ok { + t.Fatal("Expected reservation 5 to succeed") + } + dn.capacityReservations.Lock() + if r, exists := dn.capacityReservations.reservations[resId5]; exists { + r.createdAt = time.Now().Add(-2 * time.Minute) + } + dn.capacityReservations.Unlock() + + option := &VolumeGrowOption{DiskType: diskType} + if available := dn.AvailableSpaceForReservation(option); available != 5 { + t.Errorf("Expected expired reservation to free capacity in AvailableSpaceForReservation, got %d", available) + } +}