diff --git a/weed/admin/dash/config_toml.go b/weed/admin/dash/config_toml.go index c4c12a853..92532963e 100644 --- a/weed/admin/dash/config_toml.go +++ b/weed/admin/dash/config_toml.go @@ -98,6 +98,10 @@ func (cp *ConfigPersistence) ApplyMaintenanceConfigFromToml(v TomlConfig) error ecConf.ReplicaPlacement = v.GetString(k) ecChanged = true } + if k := "maintenance.erasure_coding.strict_placement"; v.IsSet(k) { + ecConf.StrictPlacement = v.GetBool(k) + ecChanged = true + } if !maintenanceChanged && !vacuumChanged && !balanceChanged && !ecChanged { return nil @@ -224,6 +228,7 @@ var pluginConfigSections = []pluginConfigSection{ "min_size_mb": int64Value, "preferred_tags": stringListValue, "replica_placement": stringValue, + "strict_placement": boolValue, }, // workers read collection_filter from the admin values, not the worker values adminKeys: map[string]func(v TomlConfig, key string) *plugin_pb.ConfigValue{ @@ -240,6 +245,10 @@ func int64Value(v TomlConfig, key string) *plugin_pb.ConfigValue { return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(v.GetInt(key))}} } +func boolValue(v TomlConfig, key string) *plugin_pb.ConfigValue { + return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: v.GetBool(key)}} +} + func stringValue(v TomlConfig, key string) *plugin_pb.ConfigValue { return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_StringValue{StringValue: v.GetString(key)}} } diff --git a/weed/command/scaffold/admin.toml b/weed/command/scaffold/admin.toml index 8f947baf6..e1af67d7a 100644 --- a/weed/command/scaffold/admin.toml +++ b/weed/command/scaffold/admin.toml @@ -64,6 +64,8 @@ # preferred_tags = ["fast", "ssd"] # EC shard placement constraint, e.g. "020"; empty uses the master default replication # replica_placement = "" +# fail planning when the placement constraints cannot be met instead of relaxing them +# strict_placement = false # max retry attempts for a failed erasure coding job # retry_limit = 1 # seconds to wait between retry attempts diff --git a/weed/pb/worker.proto b/weed/pb/worker.proto index 3ab5ce960..fa06886d0 100644 --- a/weed/pb/worker.proto +++ b/weed/pb/worker.proto @@ -380,6 +380,7 @@ message ErasureCodingTaskConfig { string collection_filter = 4; // Only process volumes from specific collections repeated string preferred_tags = 5; // Disk tags to prioritize for EC shard placement string replica_placement = 6; // EC shard replica placement (e.g. "020"); empty falls back to master default replication + bool strict_placement = 7; // fail planning instead of relaxing placement constraints } // BalanceTaskConfig contains balance-specific configuration diff --git a/weed/pb/worker_pb/worker.pb.go b/weed/pb/worker_pb/worker.pb.go index 46f4f088e..a516babcc 100644 --- a/weed/pb/worker_pb/worker.pb.go +++ b/weed/pb/worker_pb/worker.pb.go @@ -2980,6 +2980,7 @@ type ErasureCodingTaskConfig struct { CollectionFilter string `protobuf:"bytes,4,opt,name=collection_filter,json=collectionFilter,proto3" json:"collection_filter,omitempty"` // Only process volumes from specific collections PreferredTags []string `protobuf:"bytes,5,rep,name=preferred_tags,json=preferredTags,proto3" json:"preferred_tags,omitempty"` // Disk tags to prioritize for EC shard placement ReplicaPlacement string `protobuf:"bytes,6,opt,name=replica_placement,json=replicaPlacement,proto3" json:"replica_placement,omitempty"` // EC shard replica placement (e.g. "020"); empty falls back to master default replication + StrictPlacement bool `protobuf:"varint,7,opt,name=strict_placement,json=strictPlacement,proto3" json:"strict_placement,omitempty"` // fail planning instead of relaxing placement constraints unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -3056,6 +3057,13 @@ func (x *ErasureCodingTaskConfig) GetReplicaPlacement() string { return "" } +func (x *ErasureCodingTaskConfig) GetStrictPlacement() bool { + if x != nil { + return x.StrictPlacement + } + return false +} + // BalanceTaskConfig contains balance-specific configuration type BalanceTaskConfig struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -4285,14 +4293,15 @@ const file_worker_proto_rawDesc = "" + "\x10VacuumTaskConfig\x12+\n" + "\x11garbage_threshold\x18\x01 \x01(\x01R\x10garbageThreshold\x12/\n" + "\x14min_volume_age_hours\x18\x02 \x01(\x05R\x11minVolumeAgeHours\x120\n" + - "\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\x9a\x02\n" + + "\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\xc5\x02\n" + "\x17ErasureCodingTaskConfig\x12%\n" + "\x0efullness_ratio\x18\x01 \x01(\x01R\rfullnessRatio\x12*\n" + "\x11quiet_for_seconds\x18\x02 \x01(\x05R\x0fquietForSeconds\x12+\n" + "\x12min_volume_size_mb\x18\x03 \x01(\x05R\x0fminVolumeSizeMb\x12+\n" + "\x11collection_filter\x18\x04 \x01(\tR\x10collectionFilter\x12%\n" + "\x0epreferred_tags\x18\x05 \x03(\tR\rpreferredTags\x12+\n" + - "\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\"\x9b\x01\n" + + "\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\x12)\n" + + "\x10strict_placement\x18\a \x01(\bR\x0fstrictPlacement\"\x9b\x01\n" + "\x11BalanceTaskConfig\x12/\n" + "\x13imbalance_threshold\x18\x01 \x01(\x01R\x12imbalanceThreshold\x12(\n" + "\x10min_server_count\x18\x02 \x01(\x05R\x0eminServerCount\x12+\n" + diff --git a/weed/plugin/worker/config.go b/weed/plugin/worker/config.go index 0f2d66a2e..9f5ec5cda 100644 --- a/weed/plugin/worker/config.go +++ b/weed/plugin/worker/config.go @@ -125,6 +125,21 @@ func ReadIntConfig(values map[string]*plugin_pb.ConfigValue, field string, fallb return int(v) } +// ReadBoolConfig reads a bool-valued plugin config field. +func ReadBoolConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback bool) bool { + if values == nil { + return fallback + } + value := values[field] + if value == nil { + return fallback + } + if kind, ok := value.Kind.(*plugin_pb.ConfigValue_BoolValue); ok { + return kind.BoolValue + } + return fallback +} + // ReadBytesConfig reads a bytes-valued plugin config field, returning nil when // the value is missing or of a different kind. func ReadBytesConfig(values map[string]*plugin_pb.ConfigValue, field string) []byte { diff --git a/weed/shell/command_ec_encode.go b/weed/shell/command_ec_encode.go index c8943b2dd..7ae07cfcf 100644 --- a/weed/shell/command_ec_encode.go +++ b/weed/shell/command_ec_encode.go @@ -47,6 +47,15 @@ func (c *commandEcEncode) Help() string { If you only have less than 4 volume servers, with erasure coding, at least you can afford to have 4 corrupted shard files. + The guarantee follows from where shards land: a volume survives the loss of any + nodes (or racks) that hold at most parityShards (4) shards between them. Spread + is best-effort; -shardReplicaPlacement requests limits: its rack digit + sets the requested shards per rack, and its node digit the requested shards + per node (the data-center digit is not used for EC). For example + -shardReplicaPlacement=021 requests at most 1 shard per node and 2 per rack. + Only when the final placement meets these limits does losing one rack cost + at most 2 shards; the command does not guarantee that the limits are met. + The -collection parameter is a comma-separated list of collection names, with "*" and "?" wildcards, and regex patterns: - One collection: ec.encode -collection="mybucket" diff --git a/weed/worker/tasks/erasure_coding/config.go b/weed/worker/tasks/erasure_coding/config.go index 5f7c085fb..b356a61b0 100644 --- a/weed/worker/tasks/erasure_coding/config.go +++ b/weed/worker/tasks/erasure_coding/config.go @@ -18,6 +18,7 @@ type Config struct { MinSizeMB int `json:"min_size_mb"` PreferredTags []string `json:"preferred_tags"` ReplicaPlacement string `json:"replica_placement"` // e.g. "020"; empty falls back to the master default replication + StrictPlacement bool `json:"strict_placement"` // fail planning instead of relaxing placement constraints } // NewDefaultConfig creates a new default erasure coding configuration @@ -171,6 +172,18 @@ func GetConfigSpec() base.ConfigSpec { InputType: "text", CSSClasses: "form-control", }, + { + Name: "strict_placement", + JSONName: "strict_placement", + Type: config.FieldTypeBool, + DefaultValue: false, + Required: false, + DisplayName: "Strict Placement", + Description: "Refuse to encode a volume when the placement constraints can't be satisfied", + HelpText: "When enabled, a volume is only encoded if every shard can be placed within the per-disk, anti-affinity, replica-placement and per-rack caps, so the configured resilience is preserved. When disabled (default), unsatisfiable constraints are relaxed and noted in the log", + InputType: "checkbox", + CSSClasses: "form-check-input", + }, }, } } @@ -192,6 +205,7 @@ func (c *Config) ToTaskPolicy() *worker_pb.TaskPolicy { CollectionFilter: c.CollectionFilter, PreferredTags: preferredTagsCopy, ReplicaPlacement: c.ReplicaPlacement, + StrictPlacement: c.StrictPlacement, }, }, } @@ -216,6 +230,7 @@ func (c *Config) FromTaskPolicy(policy *worker_pb.TaskPolicy) error { c.CollectionFilter = ecConfig.CollectionFilter c.PreferredTags = append([]string(nil), ecConfig.PreferredTags...) c.ReplicaPlacement = ecConfig.ReplicaPlacement + c.StrictPlacement = ecConfig.StrictPlacement } return nil diff --git a/weed/worker/tasks/erasure_coding/detection.go b/weed/worker/tasks/erasure_coding/detection.go index 60234e7d5..de48d6c55 100644 --- a/weed/worker/tasks/erasure_coding/detection.go +++ b/weed/worker/tasks/erasure_coding/detection.go @@ -487,8 +487,11 @@ func buildNodeAddressMap(at *topology.ActiveTopology) map[string]string { // // Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP and, // last, the total-shards-per-rack cap as needed, failing only when no eligible -// disk has room. It prefers the source disk type but spills if that type can't -// hold every shard. rp is the resolved replica placement (may be nil). +// disk has room. ecConfig.StrictPlacement instead fails the volume's planning +// when the constraints cannot be satisfied (PlaceStrict), so a volume is only +// encoded while its configured resilience is preserved. The task prefers the +// source disk type but spills if that type can't hold every shard. rp is the +// resolved replica placement (may be nil). func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]string, metric *types.VolumeHealthMetrics, ecConfig *Config, rp *super_block.ReplicaPlacement, dataShards, parityShards int) (*topology.MultiDestinationPlan, [][]uint32, error) { if snap == nil { return nil, nil, fmt.Errorf("EC placement snapshot not available") @@ -511,13 +514,17 @@ func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]stri for i := range need { need[i] = i } + mode := ecbalancer.PlaceDurabilityFirst + if ecConfig.StrictPlacement { + mode = ecbalancer.PlaceStrict + } res, err := snap.Place(metric.VolumeID, metric.Collection, need, ecbalancer.Constraints{ DiskType: metric.DiskType, DiskTypePolicy: ecbalancer.DiskTypePrefer, PreferredTags: ecConfig.PreferredTags, ReplicaPlacement: rp, Ratio: func(string) (int, int) { return dataShards, parityShards }, - }, ecbalancer.PlaceDurabilityFirst) + }, mode) if err != nil { return nil, nil, err } diff --git a/weed/worker/tasks/erasure_coding/detection_strict_test.go b/weed/worker/tasks/erasure_coding/detection_strict_test.go new file mode 100644 index 000000000..7b7cddbb6 --- /dev/null +++ b/weed/worker/tasks/erasure_coding/detection_strict_test.go @@ -0,0 +1,50 @@ +package erasure_coding + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding/ecbalancer" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/worker/types" + "github.com/stretchr/testify/require" +) + +// StrictPlacement refuses a volume whose configured caps cannot be met, +// rather than relaxing them: here one rack must hold all 14 shards under +// a 2-shards-per-rack replica placement. +func TestPlanECDestinationsStrictPlacement(t *testing.T) { + activeTopology := buildActiveTopology(t, 7, []string{"hdd"}, 100, 0, "") + metric := &types.VolumeHealthMetrics{ + VolumeID: 1, + Server: "10.0.0.1:8080", + Size: 100 * 1024 * 1024, + } + rp, err := super_block.NewReplicaPlacementFromString("020") + require.NoError(t, err) + nodeAddresses := buildNodeAddressMap(activeTopology) + + snap := ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount) + cfg := NewDefaultConfig() + plan, shardsPerPlan, err := planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount) + require.NoError(t, err, "lenient placement relaxes the unsatisfiable rack cap") + requireAllShardsPlaced(t, plan, shardsPerPlan) + + snap = ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount) + cfg.StrictPlacement = true + _, _, err = planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount) + require.Error(t, err, "strict placement must refuse rather than weaken the rack cap") +} + +func TestStrictPlacementRoundTripsThroughTaskPolicy(t *testing.T) { + cfg := NewDefaultConfig() + cfg.StrictPlacement = true + + restored := NewDefaultConfig() + require.NoError(t, restored.FromTaskPolicy(cfg.ToTaskPolicy())) + require.True(t, restored.StrictPlacement, "strict placement must survive the persisted policy round trip") + + restored.StrictPlacement = false + require.NoError(t, restored.FromTaskPolicy(restored.ToTaskPolicy())) + require.False(t, restored.StrictPlacement) +} diff --git a/weed/worker/tasks/erasure_coding/plugin_handler.go b/weed/worker/tasks/erasure_coding/plugin_handler.go index 03b0a1b6e..1b8f74585 100644 --- a/weed/worker/tasks/erasure_coding/plugin_handler.go +++ b/weed/worker/tasks/erasure_coding/plugin_handler.go @@ -146,6 +146,13 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor { FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_STRING, Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TEXT, }, + { + Name: "strict_placement", + Label: "Strict Placement", + Description: "Fail EC planning when the placement constraints cannot be met instead of relaxing them.", + FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_BOOL, + Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TOGGLE, + }, }, }, }, @@ -165,6 +172,9 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor { "replica_placement": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""}, }, + "strict_placement": { + Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: false}, + }, }, }, AdminRuntimeDefaults: &plugin_pb.AdminRuntimeDefaults{ @@ -624,6 +634,8 @@ func deriveErasureCodingWorkerConfig(values map[string]*plugin_pb.ConfigValue) * taskConfig.ReplicaPlacement = strings.TrimSpace(pluginworker.ReadStringConfig(values, "replica_placement", taskConfig.ReplicaPlacement)) + taskConfig.StrictPlacement = pluginworker.ReadBoolConfig(values, "strict_placement", taskConfig.StrictPlacement) + return &erasureCodingWorkerConfig{ TaskConfig: taskConfig, }