From 8f80dac30fed2d56ca951ced745d5e1229595ca2 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Thu, 8 Oct 2026 22:05:02 +0800 Subject: [PATCH] ec: strict_placement option so encode only runs while guarantees hold (#11656) * ec: strict_placement option so encode only runs while guarantees hold Shard placement during encode was best-effort (PlaceDurabilityFirst): when the cluster could not satisfy the per-disk caps, anti-affinity, replica-placement or per-rack caps, the constraints were relaxed and the volume was encoded anyway, weaker than configured. A strict_placement option on the erasure coding task switches planning to PlaceStrict so the volume's planning fails instead, and the encode is retried when capacity allows the guarantee. Also documents the resilience rule in ec.encode help: a volume survives losing any nodes or racks holding at most parity-shards shards between them, and how -shardReplicaPlacement's rack and node digits bound that loss. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * ec: expose strict_placement through the plugin form and persisted task policy The admin UI, the admin.toml maintenance mapping, and the TaskPolicy serialization all dropped the new flag; add the bool field to ErasureCodingTaskConfig, the worker config form, and both conversion directions. * shell: describe shardReplicaPlacement as requested limits, not guarantees ec.encode places shards best-effort, so the configured rack/node caps only bound shard loss when the final placement actually satisfies them. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- weed/admin/dash/config_toml.go | 9 ++++ weed/command/scaffold/admin.toml | 2 + weed/pb/worker.proto | 1 + weed/pb/worker_pb/worker.pb.go | 13 ++++- weed/plugin/worker/config.go | 15 ++++++ weed/shell/command_ec_encode.go | 9 ++++ weed/worker/tasks/erasure_coding/config.go | 15 ++++++ weed/worker/tasks/erasure_coding/detection.go | 13 +++-- .../erasure_coding/detection_strict_test.go | 50 +++++++++++++++++++ .../tasks/erasure_coding/plugin_handler.go | 12 +++++ 10 files changed, 134 insertions(+), 5 deletions(-) create mode 100644 weed/worker/tasks/erasure_coding/detection_strict_test.go 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, }