From 1f6f473995a74871391344c2aaaa3a3fb7f55d20 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Sat, 2 May 2026 18:03:13 -0700 Subject: [PATCH] refactor(worker): co-locate plugin handlers with their task packages (#9301) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * refactor(worker): co-locate plugin handlers with their task packages Move every per-task plugin handler from weed/plugin/worker/ into the matching weed/worker/tasks// package, so each task owns its detection, scheduling, execution, and plugin handler in one place. Step 0 (within pluginworker, no behavior change): extract shared helpers that previously lived inside individual handler files into dedicated files and export the ones now consumed across packages. - activity.go: BuildExecutorActivity, BuildDetectorActivity - config.go: ReadStringConfig/Double/Int64/Bytes/StringList, MapTaskPriority - interval.go: ShouldSkipDetectionByInterval - volume_state.go: VolumeState + consts, FilterMetricsByVolumeState/Location - collection_filter.go: CollectionFilterMode + consts - volume_metrics.go: export CollectVolumeMetricsFromMasters, MasterAddressCandidates, FetchVolumeList - testing_senders_test.go: shared test stubs Phase 1: move the per-task plugin handlers (and the iceberg subpackage) into their task packages. weed/plugin/worker/vacuum_handler.go -> weed/worker/tasks/vacuum/plugin_handler.go weed/plugin/worker/ec_balance_handler.go -> weed/worker/tasks/ec_balance/plugin_handler.go weed/plugin/worker/erasure_coding_handler.go -> weed/worker/tasks/erasure_coding/plugin_handler.go weed/plugin/worker/volume_balance_handler.go -> weed/worker/tasks/balance/plugin_handler.go weed/plugin/worker/iceberg/ -> weed/worker/tasks/iceberg/ weed/plugin/worker/handlers/handlers.go now blank-imports all five task subpackages so their init() registrations fire. weed/command/mini.go and the worker tests construct the handler with vacuum.DefaultMaxExecutionConcurrency (the constant moved with the vacuum handler). admin_script remains in weed/plugin/worker/ because there is no underlying weed/worker/tasks/admin_script/ package to merge with. * refactor(worker): update test/plugin_workers imports for moved handlers Three handler constructors moved out of pluginworker into their task packages — update the integration test files in test/plugin_workers/ to import from the new locations: pluginworker.NewVacuumHandler -> vacuum.NewVacuumHandler pluginworker.NewVolumeBalanceHandler -> balance.NewVolumeBalanceHandler pluginworker.NewErasureCodingHandler -> erasure_coding.NewErasureCodingHandler The pluginworker import is kept where the file still uses pluginworker.WorkerOptions / pluginworker.JobHandler. * refactor(worker): update test/s3tables iceberg import path The iceberg subpackage moved from weed/plugin/worker/iceberg/ to weed/worker/tasks/iceberg/. test/s3tables/maintenance/maintenance_integration_test.go still imported the old path, breaking S3 Tables / RisingWave / Trino / Spark / Iceberg-catalog / STS integration test builds. Mirrors the OSS-side fix needed by every job in the run that transitively imports test/s3tables/maintenance. * chore: gofmt PR-touched files The S3 Tables Format Check job runs `gofmt -l` over weed/s3api/s3tables and test/s3tables, then fails if anything is unformatted. Files this PR moved or modified had import-grouping and trailing-spacing issues introduced by perl-based renames; reformat them with gofmt -w. Touched files: test/plugin_workers/erasure_coding/{detection,execution}_test.go test/s3tables/maintenance/maintenance_integration_test.go weed/plugin/worker/handlers/handlers.go weed/worker/tasks/{balance,ec_balance,erasure_coding,vacuum}/plugin_handler*.go * refactor(worker): bounds-checked int conversions for plugin config values CodeQL flagged 18 go/incorrect-integer-conversion warnings on the moved plugin handler files: results of pluginworker.ReadInt64Config (which ultimately calls strconv.ParseInt with bit size 64) were being narrowed to int32/uint32/int without an upper-bound check, so a malicious or malformed admin/worker config value could overflow the target type. Add three helpers in weed/plugin/worker/config.go that wrap ReadInt64Config and clamp out-of-range values back to the caller's fallback: ReadInt32Config (math.MinInt32 .. math.MaxInt32) ReadUint32Config (0 .. math.MaxUint32) ReadIntConfig (math.MinInt32 .. math.MaxInt32, platform-portable) Update each flagged call site in the four moved task packages to use the bounds-checked helper. For protobuf uint32 fields (volume IDs) the variable type also becomes uint32, removing the trailing uint32(volumeID) casts and changing the "missing volume_id" check from `<= 0` to `== 0`. Touched files: weed/plugin/worker/config.go weed/worker/tasks/balance/plugin_handler.go weed/worker/tasks/erasure_coding/plugin_handler.go weed/worker/tasks/vacuum/plugin_handler.go * refactor(worker): use ReadIntConfig for clamped derive-worker-config helpers CodeQL still flagged three call sites where ReadInt64Config was being narrowed to int after a value-range clamp (max_concurrent_moves <= 50, batch_size <= 100, min_server_count >= 2). The clamp is correct but CodeQL's flow analysis didn't recognize the bound, so it flagged them as unbounded narrowing. Switch to ReadIntConfig (already int32-bounded by the helper) for those three sites, drop the now-redundant int64 intermediate variables. Also drops the now-unused `> math.MaxInt32` clamp in ec_balance.deriveECBalanceWorkerConfig (the helper covers it). --- .../erasure_coding/detection_test.go | 3 +- .../erasure_coding/execution_test.go | 3 +- .../erasure_coding/large_topology_test.go | 3 +- test/plugin_workers/vacuum/detection_test.go | 3 +- test/plugin_workers/vacuum/execution_test.go | 3 +- .../volume_balance/detection_test.go | 3 +- .../volume_balance/execution_test.go | 5 +- .../maintenance_integration_test.go | 2 +- weed/command/mini.go | 6 +- weed/command/mini_plugin_test.go | 4 +- weed/command/plugin_worker_test.go | 11 +- weed/command/worker_test.go | 4 +- weed/plugin/worker/activity.go | 27 +++ weed/plugin/worker/admin_script_handler.go | 14 +- .../worker/admin_script_handler_test.go | 4 +- weed/plugin/worker/collection_filter.go | 14 ++ weed/plugin/worker/config.go | 227 ++++++++++++++++++ weed/plugin/worker/handlers/handlers.go | 6 +- weed/plugin/worker/interval.go | 20 ++ weed/plugin/worker/testing_senders_test.go | 39 +++ weed/plugin/worker/volume_metrics.go | 23 +- weed/plugin/worker/volume_metrics_test.go | 4 +- weed/plugin/worker/volume_state.go | 76 ++++++ .../tasks/balance/plugin_handler.go} | 213 ++++++---------- .../tasks/balance/plugin_handler_test.go} | 60 ++++- .../tasks/ec_balance/plugin_handler.go} | 77 +++--- .../tasks/ec_balance/plugin_handler_test.go} | 7 +- .../tasks/erasure_coding/plugin_handler.go} | 147 ++++-------- .../erasure_coding/plugin_handler_test.go} | 56 ++++- .../tasks}/iceberg/compact.go | 0 .../worker => worker/tasks}/iceberg/config.go | 0 .../tasks}/iceberg/delete_rewrite.go | 0 .../tasks}/iceberg/detection.go | 0 .../tasks}/iceberg/exec_test.go | 0 .../tasks}/iceberg/filer_io.go | 0 .../tasks}/iceberg/handler.go | 0 .../tasks}/iceberg/handler_test.go | 0 .../tasks}/iceberg/operations.go | 0 .../tasks}/iceberg/planning_index.go | 0 .../tasks}/iceberg/platform_test.go | 0 .../tasks}/iceberg/resource_groups.go | 0 .../tasks}/iceberg/sort_strategy.go | 0 .../tasks}/iceberg/testing_api.go | 0 .../tasks}/iceberg/where_filter.go | 0 .../tasks}/iceberg/where_filter_test.go | 0 .../tasks/vacuum/plugin_handler.go} | 220 +++-------------- .../tasks/vacuum/plugin_handler_test.go} | 30 +-- 47 files changed, 768 insertions(+), 546 deletions(-) create mode 100644 weed/plugin/worker/activity.go create mode 100644 weed/plugin/worker/collection_filter.go create mode 100644 weed/plugin/worker/config.go create mode 100644 weed/plugin/worker/interval.go create mode 100644 weed/plugin/worker/testing_senders_test.go create mode 100644 weed/plugin/worker/volume_state.go rename weed/{plugin/worker/volume_balance_handler.go => worker/tasks/balance/plugin_handler.go} (86%) rename weed/{plugin/worker/volume_balance_handler_test.go => worker/tasks/balance/plugin_handler_test.go} (91%) rename weed/{plugin/worker/ec_balance_handler.go => worker/tasks/ec_balance/plugin_handler.go} (88%) rename weed/{plugin/worker/ec_balance_handler_test.go => worker/tasks/ec_balance/plugin_handler_test.go} (98%) rename weed/{plugin/worker/erasure_coding_handler.go => worker/tasks/erasure_coding/plugin_handler.go} (85%) rename weed/{plugin/worker/erasure_coding_handler_test.go => worker/tasks/erasure_coding/plugin_handler_test.go} (85%) rename weed/{plugin/worker => worker/tasks}/iceberg/compact.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/config.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/delete_rewrite.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/detection.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/exec_test.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/filer_io.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/handler.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/handler_test.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/operations.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/planning_index.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/platform_test.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/resource_groups.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/sort_strategy.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/testing_api.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/where_filter.go (100%) rename weed/{plugin/worker => worker/tasks}/iceberg/where_filter_test.go (100%) rename weed/{plugin/worker/vacuum_handler.go => worker/tasks/vacuum/plugin_handler.go} (73%) rename weed/{plugin/worker/vacuum_handler_test.go => worker/tasks/vacuum/plugin_handler_test.go} (91%) diff --git a/test/plugin_workers/erasure_coding/detection_test.go b/test/plugin_workers/erasure_coding/detection_test.go index 3008e85bf..2fe27eda3 100644 --- a/test/plugin_workers/erasure_coding/detection_test.go +++ b/test/plugin_workers/erasure_coding/detection_test.go @@ -12,6 +12,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" ecstorage "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -153,7 +154,7 @@ func TestErasureCodingDetectionAcrossTopologies(t *testing.T) { master := pluginworkers.NewMasterServer(t, response) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewErasureCodingHandler(dialOption, t.TempDir()) + handler := erasure_coding.NewErasureCodingHandler(dialOption, t.TempDir()) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/erasure_coding/execution_test.go b/test/plugin_workers/erasure_coding/execution_test.go index 897950602..4201be975 100644 --- a/test/plugin_workers/erasure_coding/execution_test.go +++ b/test/plugin_workers/erasure_coding/execution_test.go @@ -12,6 +12,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" ecstorage "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -22,7 +23,7 @@ func TestErasureCodingExecutionEncodesShards(t *testing.T) { datSize := 1 * 1024 * 1024 dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewErasureCodingHandler(dialOption, t.TempDir()) + handler := erasure_coding.NewErasureCodingHandler(dialOption, t.TempDir()) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/erasure_coding/large_topology_test.go b/test/plugin_workers/erasure_coding/large_topology_test.go index 02334cc5c..504602870 100644 --- a/test/plugin_workers/erasure_coding/large_topology_test.go +++ b/test/plugin_workers/erasure_coding/large_topology_test.go @@ -10,6 +10,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -94,7 +95,7 @@ func TestErasureCodingDetectionLargeTopology(t *testing.T) { master := pluginworkers.NewMasterServer(t, response) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewErasureCodingHandler(dialOption, t.TempDir()) + handler := erasure_coding.NewErasureCodingHandler(dialOption, t.TempDir()) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/vacuum/detection_test.go b/test/plugin_workers/vacuum/detection_test.go index cbcf27502..cfcfc9630 100644 --- a/test/plugin_workers/vacuum/detection_test.go +++ b/test/plugin_workers/vacuum/detection_test.go @@ -10,6 +10,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -24,7 +25,7 @@ func TestVacuumDetectionIntegration(t *testing.T) { master := pluginworkers.NewMasterServer(t, response) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewVacuumHandler(dialOption, 1) + handler := vacuum.NewVacuumHandler(dialOption, 1) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/vacuum/execution_test.go b/test/plugin_workers/vacuum/execution_test.go index 7d749e5a4..e1846b0a9 100644 --- a/test/plugin_workers/vacuum/execution_test.go +++ b/test/plugin_workers/vacuum/execution_test.go @@ -9,6 +9,7 @@ import ( pluginworkers "github.com/seaweedfs/seaweedfs/test/plugin_workers" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -18,7 +19,7 @@ func TestVacuumExecutionIntegration(t *testing.T) { volumeID := uint32(202) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewVacuumHandler(dialOption, 1) + handler := vacuum.NewVacuumHandler(dialOption, 1) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/volume_balance/detection_test.go b/test/plugin_workers/volume_balance/detection_test.go index dd3b3dd9f..2e6fffa28 100644 --- a/test/plugin_workers/volume_balance/detection_test.go +++ b/test/plugin_workers/volume_balance/detection_test.go @@ -10,6 +10,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -21,7 +22,7 @@ func TestVolumeBalanceDetectionIntegration(t *testing.T) { master := pluginworkers.NewMasterServer(t, response) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewVolumeBalanceHandler(dialOption) + handler := balance.NewVolumeBalanceHandler(dialOption) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/plugin_workers/volume_balance/execution_test.go b/test/plugin_workers/volume_balance/execution_test.go index da5c41286..70709c6b3 100644 --- a/test/plugin_workers/volume_balance/execution_test.go +++ b/test/plugin_workers/volume_balance/execution_test.go @@ -10,6 +10,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" "github.com/stretchr/testify/require" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -22,7 +23,7 @@ func TestVolumeBalanceExecutionIntegration(t *testing.T) { volumeID := uint32(303) dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewVolumeBalanceHandler(dialOption) + handler := balance.NewVolumeBalanceHandler(dialOption) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, @@ -72,7 +73,7 @@ func TestVolumeBalanceExecutionIntegration(t *testing.T) { func TestVolumeBalanceBatchExecutionIntegration(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handler := pluginworker.NewVolumeBalanceHandler(dialOption) + handler := balance.NewVolumeBalanceHandler(dialOption) harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{ WorkerOptions: pluginworker.WorkerOptions{ GrpcDialOption: dialOption, diff --git a/test/s3tables/maintenance/maintenance_integration_test.go b/test/s3tables/maintenance/maintenance_integration_test.go index f57c2cc4e..0fd142516 100644 --- a/test/s3tables/maintenance/maintenance_integration_test.go +++ b/test/s3tables/maintenance/maintenance_integration_test.go @@ -44,8 +44,8 @@ import ( "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" - icebergHandler "github.com/seaweedfs/seaweedfs/weed/plugin/worker/iceberg" "github.com/seaweedfs/seaweedfs/weed/s3api/s3tables" + icebergHandler "github.com/seaweedfs/seaweedfs/weed/worker/tasks/iceberg" ) // --------------------------------------------------------------------------- diff --git a/weed/command/mini.go b/weed/command/mini.go index daf1c1ac0..949b1ce73 100644 --- a/weed/command/mini.go +++ b/weed/command/mini.go @@ -22,12 +22,12 @@ import ( "github.com/seaweedfs/seaweedfs/weed/util/grace" "github.com/seaweedfs/seaweedfs/weed/util/version" "github.com/seaweedfs/seaweedfs/weed/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "github.com/seaweedfs/seaweedfs/weed/worker/types" // Import task packages to trigger their auto-registration _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" - _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" ) type MiniOptions struct { @@ -1373,7 +1373,7 @@ func startMiniPluginWorker(ctx context.Context, workerDir string) { util.LoadConfiguration("security", false) grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.worker") - handlers, err := buildPluginWorkerHandlers(defaultMiniPluginJobTypes, grpcDialOption, int(pluginworker.DefaultMaxExecutionConcurrency), workerDir) + handlers, err := buildPluginWorkerHandlers(defaultMiniPluginJobTypes, grpcDialOption, int(vacuum.DefaultMaxExecutionConcurrency), workerDir) if err != nil { glog.Fatalf("Failed to build mini plugin worker handlers: %v", err) } @@ -1391,7 +1391,7 @@ func startMiniPluginWorker(ctx context.Context, workerDir string) { HeartbeatInterval: 15 * time.Second, ReconnectDelay: 5 * time.Second, MaxDetectionConcurrency: 1, - MaxExecutionConcurrency: int(pluginworker.DefaultMaxExecutionConcurrency), + MaxExecutionConcurrency: int(vacuum.DefaultMaxExecutionConcurrency), GrpcDialOption: grpcDialOption, Handlers: handlers, }) diff --git a/weed/command/mini_plugin_test.go b/weed/command/mini_plugin_test.go index 37fe694a3..36e927fa2 100644 --- a/weed/command/mini_plugin_test.go +++ b/weed/command/mini_plugin_test.go @@ -3,7 +3,7 @@ package command import ( "testing" - pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" ) @@ -11,7 +11,7 @@ import ( func TestMiniDefaultPluginJobTypes(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) // defaultMiniPluginJobTypes is "all", which includes every registered handler - handlers, err := buildPluginWorkerHandlers(defaultMiniPluginJobTypes, dialOption, int(pluginworker.DefaultMaxExecutionConcurrency), "") + handlers, err := buildPluginWorkerHandlers(defaultMiniPluginJobTypes, dialOption, int(vacuum.DefaultMaxExecutionConcurrency), "") if err != nil { t.Fatalf("buildPluginWorkerHandlers(mini default) err = %v", err) } diff --git a/weed/command/plugin_worker_test.go b/weed/command/plugin_worker_test.go index 0a3722a4b..9e90f81d5 100644 --- a/weed/command/plugin_worker_test.go +++ b/weed/command/plugin_worker_test.go @@ -10,13 +10,14 @@ import ( "testing" pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" ) func TestBuildPluginWorkerHandlerExplicitTypes(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - testMaxConcurrency := int(pluginworker.DefaultMaxExecutionConcurrency) + testMaxConcurrency := int(vacuum.DefaultMaxExecutionConcurrency) for _, jobType := range []string{"vacuum", "volume_balance", "erasure_coding", "admin_script", "iceberg_maintenance"} { handlers, err := buildPluginWorkerHandlers(jobType, dialOption, testMaxConcurrency, "") @@ -31,7 +32,7 @@ func TestBuildPluginWorkerHandlerExplicitTypes(t *testing.T) { func TestBuildPluginWorkerHandlerAliases(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - testMaxConcurrency := int(pluginworker.DefaultMaxExecutionConcurrency) + testMaxConcurrency := int(vacuum.DefaultMaxExecutionConcurrency) for _, alias := range []string{"balance", "ec", "iceberg", "admin", "script"} { handlers, err := buildPluginWorkerHandlers(alias, dialOption, testMaxConcurrency, "") @@ -54,7 +55,7 @@ func TestBuildPluginWorkerHandlerUnknown(t *testing.T) { func TestBuildPluginWorkerHandlers(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - testMaxConcurrency := int(pluginworker.DefaultMaxExecutionConcurrency) + testMaxConcurrency := int(vacuum.DefaultMaxExecutionConcurrency) handlers, err := buildPluginWorkerHandlers("vacuum,volume_balance,erasure_coding", dialOption, testMaxConcurrency, "") if err != nil { @@ -80,7 +81,7 @@ func TestBuildPluginWorkerHandlers(t *testing.T) { func TestBuildPluginWorkerHandlersCategories(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - testMaxConcurrency := int(pluginworker.DefaultMaxExecutionConcurrency) + testMaxConcurrency := int(vacuum.DefaultMaxExecutionConcurrency) allHandlers, err := buildPluginWorkerHandlers("all", dialOption, testMaxConcurrency, "") if err != nil { @@ -146,7 +147,7 @@ func TestBuildPluginWorkerHandlersCategories(t *testing.T) { func TestPluginWorkerDefaultJobTypes(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - testMaxConcurrency := int(pluginworker.DefaultMaxExecutionConcurrency) + testMaxConcurrency := int(vacuum.DefaultMaxExecutionConcurrency) // defaultPluginWorkerJobTypes is "all", so it should match the "all" category exactly defaultHandlers, err := buildPluginWorkerHandlers(defaultPluginWorkerJobTypes, dialOption, testMaxConcurrency, "") diff --git a/weed/command/worker_test.go b/weed/command/worker_test.go index 881a73aec..6d96af7d4 100644 --- a/weed/command/worker_test.go +++ b/weed/command/worker_test.go @@ -3,14 +3,14 @@ package command import ( "testing" - pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" + "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" ) func TestWorkerDefaultJobTypes(t *testing.T) { dialOption := grpc.WithTransportCredentials(insecure.NewCredentials()) - handlers, err := buildPluginWorkerHandlers(*workerJobType, dialOption, int(pluginworker.DefaultMaxExecutionConcurrency), "") + handlers, err := buildPluginWorkerHandlers(*workerJobType, dialOption, int(vacuum.DefaultMaxExecutionConcurrency), "") if err != nil { t.Fatalf("buildPluginWorkerHandlers(default worker flag) err = %v", err) } diff --git a/weed/plugin/worker/activity.go b/weed/plugin/worker/activity.go new file mode 100644 index 000000000..2ffa48092 --- /dev/null +++ b/weed/plugin/worker/activity.go @@ -0,0 +1,27 @@ +package pluginworker + +import ( + "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// BuildExecutorActivity creates an executor activity event. Exported for sub-packages. +func BuildExecutorActivity(stage string, message string) *plugin_pb.ActivityEvent { + return &plugin_pb.ActivityEvent{ + Source: plugin_pb.ActivitySource_ACTIVITY_SOURCE_EXECUTOR, + Stage: stage, + Message: message, + CreatedAt: timestamppb.Now(), + } +} + +// BuildDetectorActivity creates a detector activity event. Exported for sub-packages. +func BuildDetectorActivity(stage string, message string, details map[string]*plugin_pb.ConfigValue) *plugin_pb.ActivityEvent { + return &plugin_pb.ActivityEvent{ + Source: plugin_pb.ActivitySource_ACTIVITY_SOURCE_DETECTOR, + Stage: stage, + Message: message, + Details: details, + CreatedAt: timestamppb.Now(), + } +} diff --git a/weed/plugin/worker/admin_script_handler.go b/weed/plugin/worker/admin_script_handler.go index 26f7c4a3a..5ba14019f 100644 --- a/weed/plugin/worker/admin_script_handler.go +++ b/weed/plugin/worker/admin_script_handler.go @@ -140,8 +140,8 @@ func (h *AdminScriptHandler) Detect(ctx context.Context, request *plugin_pb.RunD return fmt.Errorf("job type %q is not handled by admin_script worker", request.JobType) } - script := normalizeAdminScript(readStringConfig(request.GetAdminConfigValues(), "script", "")) - scriptName := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "script_name", "")) + script := normalizeAdminScript(ReadStringConfig(request.GetAdminConfigValues(), "script", "")) + scriptName := strings.TrimSpace(ReadStringConfig(request.GetAdminConfigValues(), "script_name", "")) runIntervalMinutes := readAdminScriptRunIntervalMinutes(request.GetAdminConfigValues()) if ShouldSkipDetectionByInterval(request.GetLastSuccessfulRun(), runIntervalMinutes*60) { _ = sender.SendActivity(BuildDetectorActivity( @@ -228,13 +228,13 @@ func (h *AdminScriptHandler) Execute(ctx context.Context, request *plugin_pb.Exe return fmt.Errorf("job type %q is not handled by admin_script worker", request.Job.JobType) } - script := normalizeAdminScript(readStringConfig(request.Job.Parameters, "script", "")) - scriptName := strings.TrimSpace(readStringConfig(request.Job.Parameters, "script_name", "")) + script := normalizeAdminScript(ReadStringConfig(request.Job.Parameters, "script", "")) + scriptName := strings.TrimSpace(ReadStringConfig(request.Job.Parameters, "script_name", "")) if script == "" { - script = normalizeAdminScript(readStringConfig(request.GetAdminConfigValues(), "script", "")) + script = normalizeAdminScript(ReadStringConfig(request.GetAdminConfigValues(), "script", "")) } if scriptName == "" { - scriptName = strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "script_name", "")) + scriptName = strings.TrimSpace(ReadStringConfig(request.GetAdminConfigValues(), "script_name", "")) } commands := parseAdminScriptCommands(script) @@ -411,7 +411,7 @@ func (h *AdminScriptHandler) Execute(ctx context.Context, request *plugin_pb.Exe } func readAdminScriptRunIntervalMinutes(values map[string]*plugin_pb.ConfigValue) int { - runIntervalMinutes := int(readInt64Config(values, "run_interval_minutes", defaultAdminScriptRunMins)) + runIntervalMinutes := int(ReadInt64Config(values, "run_interval_minutes", defaultAdminScriptRunMins)) if runIntervalMinutes <= 0 { return defaultAdminScriptRunMins } diff --git a/weed/plugin/worker/admin_script_handler_test.go b/weed/plugin/worker/admin_script_handler_test.go index 7f2ab2236..4d29725e4 100644 --- a/weed/plugin/worker/admin_script_handler_test.go +++ b/weed/plugin/worker/admin_script_handler_test.go @@ -25,11 +25,11 @@ func TestAdminScriptDescriptorDefaults(t *testing.T) { if descriptor.AdminConfigForm == nil { t.Fatalf("expected admin config form") } - runInterval := readInt64Config(descriptor.AdminConfigForm.DefaultValues, "run_interval_minutes", 0) + runInterval := ReadInt64Config(descriptor.AdminConfigForm.DefaultValues, "run_interval_minutes", 0) if runInterval != defaultAdminScriptRunMins { t.Fatalf("unexpected run_interval_minutes default: got=%d want=%d", runInterval, defaultAdminScriptRunMins) } - script := readStringConfig(descriptor.AdminConfigForm.DefaultValues, "script", "") + script := ReadStringConfig(descriptor.AdminConfigForm.DefaultValues, "script", "") if strings.TrimSpace(script) == "" { t.Fatalf("expected non-empty default script") } diff --git a/weed/plugin/worker/collection_filter.go b/weed/plugin/worker/collection_filter.go new file mode 100644 index 000000000..7721858c7 --- /dev/null +++ b/weed/plugin/worker/collection_filter.go @@ -0,0 +1,14 @@ +package pluginworker + +// CollectionFilterMode controls how collections are interpreted during +// detection. The two recognized sentinels short-circuit regex matching: +// - CollectionFilterAll: pool every collection together (default). +// - CollectionFilterEach: run detection separately per collection. +// +// Any other non-empty value is treated as a regex. +type CollectionFilterMode string + +const ( + CollectionFilterAll CollectionFilterMode = "ALL_COLLECTIONS" + CollectionFilterEach CollectionFilterMode = "EACH_COLLECTION" +) diff --git a/weed/plugin/worker/config.go b/weed/plugin/worker/config.go new file mode 100644 index 000000000..0f2d66a2e --- /dev/null +++ b/weed/plugin/worker/config.go @@ -0,0 +1,227 @@ +package pluginworker + +import ( + "fmt" + "math" + "strconv" + "strings" + + "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" + workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" +) + +// ReadStringConfig reads a string-valued plugin config field, returning fallback +// when the value is missing or of an unsupported kind. +func ReadStringConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback string) string { + if values == nil { + return fallback + } + value := values[field] + if value == nil { + return fallback + } + switch kind := value.Kind.(type) { + case *plugin_pb.ConfigValue_StringValue: + return kind.StringValue + case *plugin_pb.ConfigValue_Int64Value: + return strconv.FormatInt(kind.Int64Value, 10) + case *plugin_pb.ConfigValue_DoubleValue: + return strconv.FormatFloat(kind.DoubleValue, 'f', -1, 64) + case *plugin_pb.ConfigValue_BoolValue: + return strconv.FormatBool(kind.BoolValue) + } + return fallback +} + +// ReadDoubleConfig reads a double-valued plugin config field, returning +// fallback when the value is missing or unparseable. +func ReadDoubleConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback float64) float64 { + if values == nil { + return fallback + } + value := values[field] + if value == nil { + return fallback + } + switch kind := value.Kind.(type) { + case *plugin_pb.ConfigValue_DoubleValue: + return kind.DoubleValue + case *plugin_pb.ConfigValue_Int64Value: + return float64(kind.Int64Value) + case *plugin_pb.ConfigValue_StringValue: + parsed, err := strconv.ParseFloat(strings.TrimSpace(kind.StringValue), 64) + if err == nil { + return parsed + } + case *plugin_pb.ConfigValue_BoolValue: + if kind.BoolValue { + return 1 + } + return 0 + } + return fallback +} + +// ReadInt64Config reads an int64-valued plugin config field, returning fallback +// when the value is missing or unparseable. +func ReadInt64Config(values map[string]*plugin_pb.ConfigValue, field string, fallback int64) int64 { + if values == nil { + return fallback + } + value := values[field] + if value == nil { + return fallback + } + switch kind := value.Kind.(type) { + case *plugin_pb.ConfigValue_Int64Value: + return kind.Int64Value + case *plugin_pb.ConfigValue_DoubleValue: + return int64(kind.DoubleValue) + case *plugin_pb.ConfigValue_StringValue: + parsed, err := strconv.ParseInt(strings.TrimSpace(kind.StringValue), 10, 64) + if err == nil { + return parsed + } + case *plugin_pb.ConfigValue_BoolValue: + if kind.BoolValue { + return 1 + } + return 0 + } + return fallback +} + +// ReadInt32Config reads an int32-valued plugin config field, returning fallback +// when the value is missing or out of int32 range. Used for protobuf int32 +// fields whose admin/worker config values arrive as int64. +func ReadInt32Config(values map[string]*plugin_pb.ConfigValue, field string, fallback int32) int32 { + v := ReadInt64Config(values, field, int64(fallback)) + if v < int64(math.MinInt32) || v > int64(math.MaxInt32) { + return fallback + } + return int32(v) +} + +// ReadUint32Config reads a uint32-valued plugin config field, returning +// fallback when the value is missing, negative, or exceeds math.MaxUint32. +// Used for protobuf uint32 fields (volume IDs, shard counts, …). +func ReadUint32Config(values map[string]*plugin_pb.ConfigValue, field string, fallback uint32) uint32 { + v := ReadInt64Config(values, field, int64(fallback)) + if v < 0 || v > int64(math.MaxUint32) { + return fallback + } + return uint32(v) +} + +// ReadIntConfig reads an int-valued plugin config field, returning fallback +// when the value is missing or outside the int32 range. The int32 range is +// used as the platform-portable safe range so that the same value parses +// identically on 32-bit and 64-bit builds. +func ReadIntConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback int) int { + v := ReadInt64Config(values, field, int64(fallback)) + if v < int64(math.MinInt32) || v > int64(math.MaxInt32) { + return fallback + } + return int(v) +} + +// 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 { + if values == nil { + return nil + } + value := values[field] + if value == nil { + return nil + } + if kind, ok := value.Kind.(*plugin_pb.ConfigValue_BytesValue); ok { + return kind.BytesValue + } + return nil +} + +// ReadStringListConfig reads a list-of-strings plugin config field, returning +// nil when the value is missing. Accepts ConfigValue_StringList, +// ConfigValue_ListValue, or a comma-separated ConfigValue_StringValue. +func ReadStringListConfig(values map[string]*plugin_pb.ConfigValue, field string) []string { + if values == nil { + return nil + } + value := values[field] + if value == nil { + return nil + } + + switch kind := value.Kind.(type) { + case *plugin_pb.ConfigValue_StringList: + return normalizeStringList(kind.StringList.GetValues()) + case *plugin_pb.ConfigValue_ListValue: + out := make([]string, 0, len(kind.ListValue.GetValues())) + for _, item := range kind.ListValue.GetValues() { + itemText := readStringFromConfigValue(item) + if itemText != "" { + out = append(out, itemText) + } + } + return normalizeStringList(out) + case *plugin_pb.ConfigValue_StringValue: + return normalizeStringList(strings.Split(kind.StringValue, ",")) + } + + return nil +} + +func readStringFromConfigValue(value *plugin_pb.ConfigValue) string { + if value == nil { + return "" + } + switch kind := value.Kind.(type) { + case *plugin_pb.ConfigValue_StringValue: + return strings.TrimSpace(kind.StringValue) + case *plugin_pb.ConfigValue_Int64Value: + return fmt.Sprintf("%d", kind.Int64Value) + case *plugin_pb.ConfigValue_DoubleValue: + return fmt.Sprintf("%g", kind.DoubleValue) + case *plugin_pb.ConfigValue_BoolValue: + if kind.BoolValue { + return "true" + } + return "false" + } + return "" +} + +func normalizeStringList(values []string) []string { + normalized := make([]string, 0, len(values)) + seen := make(map[string]struct{}, len(values)) + for _, value := range values { + item := strings.TrimSpace(value) + if item == "" { + continue + } + if _, found := seen[item]; found { + continue + } + seen[item] = struct{}{} + normalized = append(normalized, item) + } + return normalized +} + +// MapTaskPriority converts a worker-task priority into the plugin protocol's +// JobPriority enum. +func MapTaskPriority(priority workertypes.TaskPriority) plugin_pb.JobPriority { + switch strings.ToLower(string(priority)) { + case "low": + return plugin_pb.JobPriority_JOB_PRIORITY_LOW + case "medium", "normal": + return plugin_pb.JobPriority_JOB_PRIORITY_NORMAL + case "high": + return plugin_pb.JobPriority_JOB_PRIORITY_HIGH + case "critical": + return plugin_pb.JobPriority_JOB_PRIORITY_CRITICAL + default: + return plugin_pb.JobPriority_JOB_PRIORITY_NORMAL + } +} diff --git a/weed/plugin/worker/handlers/handlers.go b/weed/plugin/worker/handlers/handlers.go index fe7990e5b..c100b4898 100644 --- a/weed/plugin/worker/handlers/handlers.go +++ b/weed/plugin/worker/handlers/handlers.go @@ -5,5 +5,9 @@ package handlers import ( - _ "github.com/seaweedfs/seaweedfs/weed/plugin/worker/iceberg" // register iceberg_maintenance handler + _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" // register volume_balance handler + _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/ec_balance" // register ec_balance handler + _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" // register erasure_coding handler + _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/iceberg" // register iceberg_maintenance handler + _ "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" // register vacuum handler ) diff --git a/weed/plugin/worker/interval.go b/weed/plugin/worker/interval.go new file mode 100644 index 000000000..880d8fe12 --- /dev/null +++ b/weed/plugin/worker/interval.go @@ -0,0 +1,20 @@ +package pluginworker + +import ( + "time" + + "google.golang.org/protobuf/types/known/timestamppb" +) + +// ShouldSkipDetectionByInterval returns true when less than minIntervalSeconds +// have elapsed since lastSuccessfulRun. Exported so sub-packages can reuse it. +func ShouldSkipDetectionByInterval(lastSuccessfulRun *timestamppb.Timestamp, minIntervalSeconds int) bool { + if lastSuccessfulRun == nil || minIntervalSeconds <= 0 { + return false + } + lastRun := lastSuccessfulRun.AsTime() + if lastRun.IsZero() { + return false + } + return time.Since(lastRun) < time.Duration(minIntervalSeconds)*time.Second +} diff --git a/weed/plugin/worker/testing_senders_test.go b/weed/plugin/worker/testing_senders_test.go new file mode 100644 index 000000000..7fb5ab2e6 --- /dev/null +++ b/weed/plugin/worker/testing_senders_test.go @@ -0,0 +1,39 @@ +package pluginworker + +import ( + "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" +) + +type noopDetectionSender struct{} + +func (noopDetectionSender) SendProposals(*plugin_pb.DetectionProposals) error { return nil } +func (noopDetectionSender) SendComplete(*plugin_pb.DetectionComplete) error { return nil } +func (noopDetectionSender) SendActivity(*plugin_pb.ActivityEvent) error { return nil } + +type noopExecutionSender struct{} + +func (noopExecutionSender) SendProgress(*plugin_pb.JobProgressUpdate) error { return nil } +func (noopExecutionSender) SendCompleted(*plugin_pb.JobCompleted) error { return nil } + +type recordingDetectionSender struct { + proposals *plugin_pb.DetectionProposals + complete *plugin_pb.DetectionComplete + events []*plugin_pb.ActivityEvent +} + +func (r *recordingDetectionSender) SendProposals(proposals *plugin_pb.DetectionProposals) error { + r.proposals = proposals + return nil +} + +func (r *recordingDetectionSender) SendComplete(complete *plugin_pb.DetectionComplete) error { + r.complete = complete + return nil +} + +func (r *recordingDetectionSender) SendActivity(event *plugin_pb.ActivityEvent) error { + if event != nil { + r.events = append(r.events, event) + } + return nil +} diff --git a/weed/plugin/worker/volume_metrics.go b/weed/plugin/worker/volume_metrics.go index aef4c0ff7..33957b88f 100644 --- a/weed/plugin/worker/volume_metrics.go +++ b/weed/plugin/worker/volume_metrics.go @@ -17,7 +17,10 @@ import ( "google.golang.org/grpc" ) -func collectVolumeMetricsFromMasters( +// CollectVolumeMetricsFromMasters dials the provided master addresses in order +// until one returns a usable volume list, then converts that into per-volume +// health metrics, an active-topology view, and a replica-location map. +func CollectVolumeMetricsFromMasters( ctx context.Context, masterAddresses []string, collectionFilter string, @@ -31,7 +34,7 @@ func collectVolumeMetricsFromMasters( } for _, masterAddress := range masterAddresses { - response, err := fetchVolumeList(ctx, masterAddress, grpcDialOption) + response, err := FetchVolumeList(ctx, masterAddress, grpcDialOption) if err != nil { glog.Warningf("Plugin worker failed master volume list at %s: %v", masterAddress, err) continue @@ -53,9 +56,12 @@ func collectVolumeMetricsFromMasters( return nil, nil, nil, fmt.Errorf("failed to load topology from all provided masters") } -func fetchVolumeList(ctx context.Context, address string, grpcDialOption grpc.DialOption) (*master_pb.VolumeListResponse, error) { +// FetchVolumeList dials the given master address (trying both the address as +// given and the gRPC port variant) and returns the master's volume list. Used +// by detection helpers that already know which master address to talk to. +func FetchVolumeList(ctx context.Context, address string, grpcDialOption grpc.DialOption) (*master_pb.VolumeListResponse, error) { var lastErr error - for _, candidate := range masterAddressCandidates(address) { + for _, candidate := range MasterAddressCandidates(address) { if ctx.Err() != nil { return nil, ctx.Err() } @@ -101,8 +107,8 @@ func buildVolumeMetrics( var collectionRegex *regexp.Regexp trimmedFilter := strings.TrimSpace(collectionFilter) - filterMode := collectionFilterMode(trimmedFilter) - if trimmedFilter != "" && filterMode != collectionFilterAll && filterMode != collectionFilterEach && trimmedFilter != "*" { + filterMode := CollectionFilterMode(trimmedFilter) + if trimmedFilter != "" && filterMode != CollectionFilterAll && filterMode != CollectionFilterEach && trimmedFilter != "*" { var err error collectionRegex, err = regexp.Compile(trimmedFilter) if err != nil { @@ -186,7 +192,10 @@ func isConfigError(err error) bool { return errors.As(err, &ce) } -func masterAddressCandidates(address string) []string { +// MasterAddressCandidates returns address forms to try when dialing a master: +// the address as given plus the gRPC variant (port + 10000). Both are tried +// because callers may pass either an HTTP-port or gRPC-port address. +func MasterAddressCandidates(address string) []string { trimmed := strings.TrimSpace(address) if trimmed == "" { return nil diff --git a/weed/plugin/worker/volume_metrics_test.go b/weed/plugin/worker/volume_metrics_test.go index 01f5ea01f..05eeb008f 100644 --- a/weed/plugin/worker/volume_metrics_test.go +++ b/weed/plugin/worker/volume_metrics_test.go @@ -53,7 +53,7 @@ func TestBuildVolumeMetricsAllCollections(t *testing.T) { &master_pb.VolumeInformationMessage{Id: 1, Collection: "photos", Size: 100}, &master_pb.VolumeInformationMessage{Id: 2, Collection: "videos", Size: 200}, ) - metrics, _, _, err := buildVolumeMetrics(resp, string(collectionFilterAll)) + metrics, _, _, err := buildVolumeMetrics(resp, string(CollectionFilterAll)) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -68,7 +68,7 @@ func TestBuildVolumeMetricsEachCollection(t *testing.T) { &master_pb.VolumeInformationMessage{Id: 2, Collection: "videos", Size: 200}, ) // EACH_COLLECTION passes all volumes through; filtering happens in the handler - metrics, _, _, err := buildVolumeMetrics(resp, string(collectionFilterEach)) + metrics, _, _, err := buildVolumeMetrics(resp, string(CollectionFilterEach)) if err != nil { t.Fatalf("unexpected error: %v", err) } diff --git a/weed/plugin/worker/volume_state.go b/weed/plugin/worker/volume_state.go new file mode 100644 index 000000000..c03f65b27 --- /dev/null +++ b/weed/plugin/worker/volume_state.go @@ -0,0 +1,76 @@ +package pluginworker + +import ( + "github.com/seaweedfs/seaweedfs/weed/util/wildcard" + workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" +) + +// VolumeState controls which volumes participate in detection (e.g. balance, +// vacuum). It is passed to FilterMetricsByVolumeState by the various plugin +// handlers. +type VolumeState string + +const ( + VolumeStateAll VolumeState = "ALL" + VolumeStateActive VolumeState = "ACTIVE" + VolumeStateFull VolumeState = "FULL" +) + +// FilterMetricsByVolumeState filters volume metrics by state. +// "ACTIVE" keeps volumes with FullnessRatio < 1.01 (writable, below size limit). +// "FULL" keeps volumes with FullnessRatio >= 1.01 (read-only, above size limit). +// "ALL" or any other value returns all metrics unfiltered. +func FilterMetricsByVolumeState(metrics []*workertypes.VolumeHealthMetrics, state VolumeState) []*workertypes.VolumeHealthMetrics { + const fullnessThreshold = 1.01 + + var predicate func(m *workertypes.VolumeHealthMetrics) bool + switch state { + case VolumeStateActive: + predicate = func(m *workertypes.VolumeHealthMetrics) bool { + return m.FullnessRatio < fullnessThreshold + } + case VolumeStateFull: + predicate = func(m *workertypes.VolumeHealthMetrics) bool { + return m.FullnessRatio >= fullnessThreshold + } + default: + return metrics + } + + filtered := make([]*workertypes.VolumeHealthMetrics, 0, len(metrics)) + for _, m := range metrics { + if m == nil { + continue + } + if predicate(m) { + filtered = append(filtered, m) + } + } + return filtered +} + +// FilterMetricsByLocation filters volume metrics by data center, rack, and node +// wildcards. Empty filters match anything in their respective dimension. +func FilterMetricsByLocation(metrics []*workertypes.VolumeHealthMetrics, dcFilter, rackFilter, nodeFilter string) []*workertypes.VolumeHealthMetrics { + dcMatchers := wildcard.CompileWildcardMatchers(dcFilter) + rackMatchers := wildcard.CompileWildcardMatchers(rackFilter) + nodeMatchers := wildcard.CompileWildcardMatchers(nodeFilter) + + filtered := make([]*workertypes.VolumeHealthMetrics, 0, len(metrics)) + for _, m := range metrics { + if m == nil { + continue + } + if !wildcard.MatchesAnyWildcard(dcMatchers, m.DataCenter) { + continue + } + if !wildcard.MatchesAnyWildcard(rackMatchers, m.Rack) { + continue + } + if !wildcard.MatchesAnyWildcard(nodeMatchers, m.Server) { + continue + } + filtered = append(filtered, m) + } + return filtered +} diff --git a/weed/plugin/worker/volume_balance_handler.go b/weed/worker/tasks/balance/plugin_handler.go similarity index 86% rename from weed/plugin/worker/volume_balance_handler.go rename to weed/worker/tasks/balance/plugin_handler.go index b12fb207c..7fb36c3fe 100644 --- a/weed/plugin/worker/volume_balance_handler.go +++ b/weed/worker/tasks/balance/plugin_handler.go @@ -1,4 +1,4 @@ -package pluginworker +package balance import ( "context" @@ -14,8 +14,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" - "github.com/seaweedfs/seaweedfs/weed/util/wildcard" - balancetask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/grpc" "google.golang.org/protobuf/proto" @@ -26,36 +25,19 @@ const ( maxProposalStringLength = 200 ) -// collectionFilterMode controls how collections are handled during balance detection. -type collectionFilterMode string - -const ( - collectionFilterAll collectionFilterMode = "ALL_COLLECTIONS" - collectionFilterEach collectionFilterMode = "EACH_COLLECTION" -) - -// volumeState controls which volumes participate in balance detection. -type volumeState string - -const ( - volumeStateAll volumeState = "ALL" - volumeStateActive volumeState = "ACTIVE" - volumeStateFull volumeState = "FULL" -) - func init() { - RegisterHandler(HandlerFactory{ + pluginworker.RegisterHandler(pluginworker.HandlerFactory{ JobType: "volume_balance", - Category: CategoryDefault, + Category: pluginworker.CategoryDefault, Aliases: []string{"balance", "volume.balance", "volume-balance"}, - Build: func(opts HandlerBuildOptions) (JobHandler, error) { + Build: func(opts pluginworker.HandlerBuildOptions) (pluginworker.JobHandler, error) { return NewVolumeBalanceHandler(opts.GrpcDialOption), nil }, }) } type volumeBalanceWorkerConfig struct { - TaskConfig *balancetask.Config + TaskConfig *Config MaxConcurrentMoves int BatchSize int } @@ -114,9 +96,9 @@ func (h *VolumeBalanceHandler) Descriptor() *plugin_pb.JobTypeDescriptor { FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_ENUM, Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_SELECT, Options: []*plugin_pb.ConfigOption{ - {Value: string(volumeStateAll), Label: "All Volumes"}, - {Value: string(volumeStateActive), Label: "Active (writable)"}, - {Value: string(volumeStateFull), Label: "Full (read-only)"}, + {Value: string(pluginworker.VolumeStateAll), Label: "All Volumes"}, + {Value: string(pluginworker.VolumeStateActive), Label: "Active (writable)"}, + {Value: string(pluginworker.VolumeStateFull), Label: "Full (read-only)"}, }, }, { @@ -151,7 +133,7 @@ func (h *VolumeBalanceHandler) Descriptor() *plugin_pb.JobTypeDescriptor { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""}, }, "volume_state": { - Kind: &plugin_pb.ConfigValue_StringValue{StringValue: string(volumeStateAll)}, + Kind: &plugin_pb.ConfigValue_StringValue{StringValue: string(pluginworker.VolumeStateAll)}, }, "data_center_filter": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""}, @@ -268,7 +250,7 @@ func (h *VolumeBalanceHandler) Descriptor() *plugin_pb.JobTypeDescriptor { func (h *VolumeBalanceHandler) Detect( ctx context.Context, request *plugin_pb.RunDetectionRequest, - sender DetectionSender, + sender pluginworker.DetectionSender, ) error { if request == nil { return fmt.Errorf("run detection request is nil") @@ -281,7 +263,7 @@ func (h *VolumeBalanceHandler) Detect( } workerConfig := deriveBalanceWorkerConfig(request.GetWorkerConfigValues()) - collectionFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) + collectionFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) masters := make([]string, 0) if request.ClusterContext != nil { masters = append(masters, request.ClusterContext.MasterGrpcAddresses...) @@ -292,15 +274,15 @@ func (h *VolumeBalanceHandler) Detect( return err } - volState := volumeState(strings.ToUpper(strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "volume_state", string(volumeStateAll))))) - metrics = filterMetricsByVolumeState(metrics, volState) + volState := pluginworker.VolumeState(strings.ToUpper(strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "volume_state", string(pluginworker.VolumeStateAll))))) + metrics = pluginworker.FilterMetricsByVolumeState(metrics, volState) - dataCenterFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) - rackFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "rack_filter", "")) - nodeFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "node_filter", "")) + dataCenterFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) + rackFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "rack_filter", "")) + nodeFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "node_filter", "")) if dataCenterFilter != "" || rackFilter != "" || nodeFilter != "" { - metrics = filterMetricsByLocation(metrics, dataCenterFilter, rackFilter, nodeFilter) + metrics = pluginworker.FilterMetricsByLocation(metrics, dataCenterFilter, rackFilter, nodeFilter) } workerConfig.TaskConfig.DataCenterFilter = dataCenterFilter @@ -316,7 +298,7 @@ func (h *VolumeBalanceHandler) Detect( var results []*workertypes.TaskDetectionResult var hasMore bool - if collectionFilterMode(collectionFilter) == collectionFilterEach { + if pluginworker.CollectionFilterMode(collectionFilter) == pluginworker.CollectionFilterEach { // Group metrics by collection in a single pass (O(N) instead of O(C*N)) metricsByCollection := make(map[string][]*workertypes.VolumeHealthMetrics) for _, m := range metrics { @@ -342,7 +324,7 @@ func (h *VolumeBalanceHandler) Detect( if unlimitedBudget { perCollectionLimit = 0 // Detection treats <= 0 as unbounded } - perResults, perHasMore, perErr := balancetask.Detection(metricsByCollection[collection], clusterInfo, workerConfig.TaskConfig, perCollectionLimit) + perResults, perHasMore, perErr := Detection(metricsByCollection[collection], clusterInfo, workerConfig.TaskConfig, perCollectionLimit) if perErr != nil { return perErr } @@ -356,7 +338,7 @@ func (h *VolumeBalanceHandler) Detect( } } else { var err error - results, hasMore, err = balancetask.Detection(metrics, clusterInfo, workerConfig.TaskConfig, maxResults) + results, hasMore, err = Detection(metrics, clusterInfo, workerConfig.TaskConfig, maxResults) if err != nil { return err } @@ -397,10 +379,10 @@ func (h *VolumeBalanceHandler) Detect( } func emitVolumeBalanceDetectionDecisionTrace( - sender DetectionSender, + sender pluginworker.DetectionSender, metrics []*workertypes.VolumeHealthMetrics, activeTopology *topology.ActiveTopology, - taskConfig *balancetask.Config, + taskConfig *Config, results []*workertypes.TaskDetectionResult, ) error { if sender == nil || taskConfig == nil { @@ -428,7 +410,7 @@ func emitVolumeBalanceDetectionDecisionTrace( ) } - if err := sender.SendActivity(BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ "total_volumes": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(totalVolumes)}, }, @@ -475,7 +457,7 @@ func emitVolumeBalanceDetectionDecisionTrace( volumeCount, minVolumeCount, ) - if err := sender.SendActivity(BuildDetectorActivity("decision_disk_type", message, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_disk_type", message, map[string]*plugin_pb.ConfigValue{ "disk_type": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: diskType}, }, @@ -496,7 +478,7 @@ func emitVolumeBalanceDetectionDecisionTrace( } // Seed server counts from topology so zero-volume servers are included, - // matching the same logic used in balancetask.Detection. + // matching the same logic used in Detection. serverVolumeCounts := make(map[string]int) if activeTopology != nil { topologyInfo := activeTopology.GetTopologyInfo() @@ -524,7 +506,7 @@ func emitVolumeBalanceDetectionDecisionTrace( len(serverVolumeCounts), taskConfig.MinServerCount, ) - if err := sender.SendActivity(BuildDetectorActivity("decision_disk_type", message, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_disk_type", message, map[string]*plugin_pb.ConfigValue{ "disk_type": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: diskType}, }, @@ -595,7 +577,7 @@ func emitVolumeBalanceDetectionDecisionTrace( ) } - if err := sender.SendActivity(BuildDetectorActivity(stage, message, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity(stage, message, map[string]*plugin_pb.ConfigValue{ "disk_type": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: diskType}, }, @@ -633,39 +615,6 @@ func emitVolumeBalanceDetectionDecisionTrace( return nil } -// filterMetricsByVolumeState filters volume metrics by state. -// "ACTIVE" keeps volumes with FullnessRatio < 1.01 (writable, below size limit). -// "FULL" keeps volumes with FullnessRatio >= 1.01 (read-only, above size limit). -// "ALL" or any other value returns all metrics unfiltered. -func filterMetricsByVolumeState(metrics []*workertypes.VolumeHealthMetrics, state volumeState) []*workertypes.VolumeHealthMetrics { - const fullnessThreshold = 1.01 - - var predicate func(m *workertypes.VolumeHealthMetrics) bool - switch state { - case volumeStateActive: - predicate = func(m *workertypes.VolumeHealthMetrics) bool { - return m.FullnessRatio < fullnessThreshold - } - case volumeStateFull: - predicate = func(m *workertypes.VolumeHealthMetrics) bool { - return m.FullnessRatio >= fullnessThreshold - } - default: - return metrics - } - - filtered := make([]*workertypes.VolumeHealthMetrics, 0, len(metrics)) - for _, m := range metrics { - if m == nil { - continue - } - if predicate(m) { - filtered = append(filtered, m) - } - } - return filtered -} - func countBalanceDiskTypes(metrics []*workertypes.VolumeHealthMetrics) int { diskTypes := make(map[string]struct{}) for _, metric := range metrics { @@ -690,7 +639,7 @@ const ( func (h *VolumeBalanceHandler) Execute( ctx context.Context, request *plugin_pb.ExecuteJobRequest, - sender ExecutionSender, + sender pluginworker.ExecutionSender, ) error { if request == nil || request.Job == nil { return fmt.Errorf("execute request/job is nil") @@ -722,7 +671,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( ctx context.Context, request *plugin_pb.ExecuteJobRequest, params *worker_pb.TaskParams, - sender ExecutionSender, + sender pluginworker.ExecutionSender, ) error { if len(params.Sources) == 0 || strings.TrimSpace(params.Sources[0].Node) == "" { return fmt.Errorf("volume balance source node is required") @@ -731,7 +680,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( return fmt.Errorf("volume balance target node is required") } - task := balancetask.NewBalanceTask( + task := NewBalanceTask( request.Job.JobId, params.Sources[0].Node, params.VolumeId, @@ -753,7 +702,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( Stage: stage, Message: message, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity(stage, message), + pluginworker.BuildExecutorActivity(stage, message), }, }); err != nil { execCancel() @@ -768,7 +717,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( Stage: "assigned", Message: "volume balance job accepted", Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("assigned", "volume balance job accepted"), + pluginworker.BuildExecutorActivity("assigned", "volume balance job accepted"), }, }); err != nil { return err @@ -783,7 +732,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( Stage: "failed", Message: err.Error(), Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("failed", err.Error()), + pluginworker.BuildExecutorActivity("failed", err.Error()), }, }) return err @@ -812,7 +761,7 @@ func (h *VolumeBalanceHandler) executeSingleMove( }, }, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("completed", resultSummary), + pluginworker.BuildExecutorActivity("completed", resultSummary), }, }) } @@ -822,7 +771,7 @@ func (h *VolumeBalanceHandler) executeBatchMoves( ctx context.Context, request *plugin_pb.ExecuteJobRequest, params *worker_pb.TaskParams, - sender ExecutionSender, + sender pluginworker.ExecutionSender, ) error { bp := params.GetBalanceParams() if len(bp.Moves) == 0 { @@ -870,7 +819,7 @@ func (h *VolumeBalanceHandler) executeBatchMoves( Stage: "assigned", Message: fmt.Sprintf("batch volume balance accepted: %d moves", totalMoves), Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("assigned", fmt.Sprintf("batch volume balance: %d moves, concurrency %d", totalMoves, maxConcurrent)), + pluginworker.BuildExecutorActivity("assigned", fmt.Sprintf("batch volume balance: %d moves, concurrency %d", totalMoves, maxConcurrent)), }, }); err != nil { return err @@ -914,7 +863,7 @@ func (h *VolumeBalanceHandler) executeBatchMoves( Stage: fmt.Sprintf("move %d/%d", moveIndex+1, totalMoves), Message: message, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity(fmt.Sprintf("move-%d", moveIndex+1), message), + pluginworker.BuildExecutorActivity(fmt.Sprintf("move-%d", moveIndex+1), message), }, }); err != nil { sendErr = err @@ -938,7 +887,7 @@ func (h *VolumeBalanceHandler) executeBatchMoves( go func(idx int, m *worker_pb.BalanceMoveSpec) { defer func() { <-sem }() // release slot - task := balancetask.NewBalanceTask( + task := NewBalanceTask( fmt.Sprintf("%s-move-%d", request.Job.JobId, idx), m.SourceNode, m.VolumeId, @@ -1010,7 +959,7 @@ func (h *VolumeBalanceHandler) executeBatchMoves( }, }, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("completed", summary), + pluginworker.BuildExecutorActivity("completed", summary), }, }) } @@ -1049,13 +998,13 @@ func (h *VolumeBalanceHandler) collectVolumeMetrics( masterAddresses []string, collectionFilter string, ) ([]*workertypes.VolumeHealthMetrics, *topology.ActiveTopology, map[uint32][]workertypes.ReplicaLocation, error) { - return collectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) + return pluginworker.CollectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) } func deriveBalanceWorkerConfig(values map[string]*plugin_pb.ConfigValue) *volumeBalanceWorkerConfig { - taskConfig := balancetask.NewDefaultConfig() + taskConfig := NewDefaultConfig() - imbalanceThreshold := readDoubleConfig(values, "imbalance_threshold", taskConfig.ImbalanceThreshold) + imbalanceThreshold := pluginworker.ReadDoubleConfig(values, "imbalance_threshold", taskConfig.ImbalanceThreshold) if imbalanceThreshold < 0 { imbalanceThreshold = 0 } @@ -1064,29 +1013,27 @@ func deriveBalanceWorkerConfig(values map[string]*plugin_pb.ConfigValue) *volume } taskConfig.ImbalanceThreshold = imbalanceThreshold - minServerCount := int(readInt64Config(values, "min_server_count", int64(taskConfig.MinServerCount))) + minServerCount := pluginworker.ReadIntConfig(values, "min_server_count", taskConfig.MinServerCount) if minServerCount < 2 { minServerCount = 2 } taskConfig.MinServerCount = minServerCount - maxConcurrentMoves64 := readInt64Config(values, "max_concurrent_moves", int64(defaultMaxConcurrentMoves)) - if maxConcurrentMoves64 < 1 { - maxConcurrentMoves64 = 1 + maxConcurrentMoves := pluginworker.ReadIntConfig(values, "max_concurrent_moves", defaultMaxConcurrentMoves) + if maxConcurrentMoves < 1 { + maxConcurrentMoves = 1 } - if maxConcurrentMoves64 > 50 { - maxConcurrentMoves64 = 50 + if maxConcurrentMoves > 50 { + maxConcurrentMoves = 50 } - maxConcurrentMoves := int(maxConcurrentMoves64) - batchSize64 := readInt64Config(values, "batch_size", 20) - if batchSize64 < 1 { - batchSize64 = 1 + batchSize := pluginworker.ReadIntConfig(values, "batch_size", 20) + if batchSize < 1 { + batchSize = 1 } - if batchSize64 > 100 { - batchSize64 = 100 + if batchSize > 100 { + batchSize = 100 } - batchSize := int(batchSize64) return &volumeBalanceWorkerConfig{ TaskConfig: taskConfig, @@ -1095,30 +1042,6 @@ func deriveBalanceWorkerConfig(values map[string]*plugin_pb.ConfigValue) *volume } } -func filterMetricsByLocation(metrics []*workertypes.VolumeHealthMetrics, dcFilter, rackFilter, nodeFilter string) []*workertypes.VolumeHealthMetrics { - dcMatchers := wildcard.CompileWildcardMatchers(dcFilter) - rackMatchers := wildcard.CompileWildcardMatchers(rackFilter) - nodeMatchers := wildcard.CompileWildcardMatchers(nodeFilter) - - filtered := make([]*workertypes.VolumeHealthMetrics, 0, len(metrics)) - for _, m := range metrics { - if m == nil { - continue - } - if !wildcard.MatchesAnyWildcard(dcMatchers, m.DataCenter) { - continue - } - if !wildcard.MatchesAnyWildcard(rackMatchers, m.Rack) { - continue - } - if !wildcard.MatchesAnyWildcard(nodeMatchers, m.Server) { - continue - } - filtered = append(filtered, m) - } - return filtered -} - func buildVolumeBalanceProposal( result *workertypes.TaskDetectionResult, ) (*plugin_pb.JobProposal, error) { @@ -1171,7 +1094,7 @@ func buildVolumeBalanceProposal( ProposalId: proposalID, DedupeKey: dedupeKey, JobType: "volume_balance", - Priority: mapTaskPriority(result.Priority), + Priority: pluginworker.MapTaskPriority(result.Priority), Summary: summary, Detail: strings.TrimSpace(result.Reason), Parameters: map[string]*plugin_pb.ConfigValue{ @@ -1371,7 +1294,7 @@ func buildBatchVolumeBalanceProposals( ProposalId: proposalID, DedupeKey: compositeDedupeKey, JobType: "volume_balance", - Priority: mapTaskPriority(highestPriority), + Priority: pluginworker.MapTaskPriority(highestPriority), Summary: summary, Detail: fmt.Sprintf("Batch of %d volume moves with concurrency %d", len(moves), maxConcurrentMoves), Parameters: map[string]*plugin_pb.ConfigValue{ @@ -1401,7 +1324,7 @@ func decodeVolumeBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return nil, fmt.Errorf("job spec is nil") } - if payload := readBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { + if payload := pluginworker.ReadBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { params := &worker_pb.TaskParams{} if err := proto.Unmarshal(payload, params); err != nil { return nil, fmt.Errorf("unmarshal task_params_pb: %w", err) @@ -1412,23 +1335,23 @@ func decodeVolumeBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return params, nil } - volumeID := readInt64Config(job.Parameters, "volume_id", 0) - sourceNode := strings.TrimSpace(readStringConfig(job.Parameters, "source_server", "")) + volumeID := pluginworker.ReadUint32Config(job.Parameters, "volume_id", 0) + sourceNode := strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "source_server", "")) if sourceNode == "" { - sourceNode = strings.TrimSpace(readStringConfig(job.Parameters, "server", "")) + sourceNode = strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "server", "")) } - targetNode := strings.TrimSpace(readStringConfig(job.Parameters, "target_server", "")) + targetNode := strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "target_server", "")) if targetNode == "" { - targetNode = strings.TrimSpace(readStringConfig(job.Parameters, "target", "")) + targetNode = strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "target", "")) } - collection := readStringConfig(job.Parameters, "collection", "") - timeoutSeconds := int32(readInt64Config(job.Parameters, "timeout_seconds", int64(defaultBalanceTimeoutSeconds))) + collection := pluginworker.ReadStringConfig(job.Parameters, "collection", "") + timeoutSeconds := pluginworker.ReadInt32Config(job.Parameters, "timeout_seconds", defaultBalanceTimeoutSeconds) if timeoutSeconds <= 0 { timeoutSeconds = defaultBalanceTimeoutSeconds } forceMove := readBoolConfig(job.Parameters, "force_move", false) - if volumeID <= 0 { + if volumeID == 0 { return nil, fmt.Errorf("missing volume_id in job parameters") } if sourceNode == "" { @@ -1440,18 +1363,18 @@ func decodeVolumeBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return &worker_pb.TaskParams{ TaskId: job.JobId, - VolumeId: uint32(volumeID), + VolumeId: volumeID, Collection: collection, Sources: []*worker_pb.TaskSource{ { Node: sourceNode, - VolumeId: uint32(volumeID), + VolumeId: volumeID, }, }, Targets: []*worker_pb.TaskTarget{ { Node: targetNode, - VolumeId: uint32(volumeID), + VolumeId: volumeID, }, }, TaskParams: &worker_pb.TaskParams_BalanceParams{ diff --git a/weed/plugin/worker/volume_balance_handler_test.go b/weed/worker/tasks/balance/plugin_handler_test.go similarity index 91% rename from weed/plugin/worker/volume_balance_handler_test.go rename to weed/worker/tasks/balance/plugin_handler_test.go index fb03db7ad..205c7fdf5 100644 --- a/weed/plugin/worker/volume_balance_handler_test.go +++ b/weed/worker/tasks/balance/plugin_handler_test.go @@ -1,4 +1,4 @@ -package pluginworker +package balance import ( "context" @@ -9,7 +9,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" - balancetask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/balance" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/protobuf/proto" ) @@ -342,7 +342,7 @@ func TestVolumeBalanceHandlerRejectsUnsupportedJobType(t *testing.T) { func TestEmitVolumeBalanceDetectionDecisionTraceNoTasks(t *testing.T) { sender := &recordingDetectionSender{} - config := balancetask.NewDefaultConfig() + config := NewDefaultConfig() config.ImbalanceThreshold = 0.2 config.MinServerCount = 2 @@ -600,31 +600,31 @@ func TestFilterMetricsByLocation(t *testing.T) { } // Filter by DC - filtered := filterMetricsByLocation(metrics, "dc1", "", "") + filtered := pluginworker.FilterMetricsByLocation(metrics, "dc1", "", "") if len(filtered) != 2 { t.Fatalf("DC filter: expected 2, got %d", len(filtered)) } // Filter by rack - filtered = filterMetricsByLocation(metrics, "", "rack1,rack2", "") + filtered = pluginworker.FilterMetricsByLocation(metrics, "", "rack1,rack2", "") if len(filtered) != 3 { t.Fatalf("rack filter: expected 3, got %d", len(filtered)) } // Filter by node - filtered = filterMetricsByLocation(metrics, "", "", "node-a,node-c") + filtered = pluginworker.FilterMetricsByLocation(metrics, "", "", "node-a,node-c") if len(filtered) != 2 { t.Fatalf("node filter: expected 2, got %d", len(filtered)) } // Combined DC + rack - filtered = filterMetricsByLocation(metrics, "dc2", "rack3", "") + filtered = pluginworker.FilterMetricsByLocation(metrics, "dc2", "rack3", "") if len(filtered) != 1 { t.Fatalf("DC+rack filter: expected 1, got %d", len(filtered)) } // Empty filters pass all - filtered = filterMetricsByLocation(metrics, "", "", "") + filtered = pluginworker.FilterMetricsByLocation(metrics, "", "", "") if len(filtered) != 4 { t.Fatalf("no filter: expected 4, got %d", len(filtered)) } @@ -674,7 +674,7 @@ func TestFilterMetricsByVolumeState(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - result := filterMetricsByVolumeState(metrics, volumeState(tt.state)) + result := pluginworker.FilterMetricsByVolumeState(metrics, pluginworker.VolumeState(tt.state)) if len(result) != len(tt.expectedIDs) { t.Fatalf("expected %d metrics, got %d", len(tt.expectedIDs), len(result)) } @@ -694,23 +694,23 @@ func TestFilterMetricsByVolumeState_NilElement(t *testing.T) { nil, {VolumeID: 2, FullnessRatio: 1.5}, } - result := filterMetricsByVolumeState(metrics, volumeStateActive) + result := pluginworker.FilterMetricsByVolumeState(metrics, pluginworker.VolumeStateActive) if len(result) != 1 || result[0].VolumeID != 1 { t.Fatalf("expected [vol 1] for ACTIVE with nil elements, got %d results", len(result)) } - result = filterMetricsByVolumeState(metrics, volumeStateFull) + result = pluginworker.FilterMetricsByVolumeState(metrics, pluginworker.VolumeStateFull) if len(result) != 1 || result[0].VolumeID != 2 { t.Fatalf("expected [vol 2] for FULL with nil elements, got %d results", len(result)) } } func TestFilterMetricsByVolumeState_EmptyInput(t *testing.T) { - result := filterMetricsByVolumeState(nil, volumeStateActive) + result := pluginworker.FilterMetricsByVolumeState(nil, pluginworker.VolumeStateActive) if len(result) != 0 { t.Fatalf("expected 0 metrics for nil input, got %d", len(result)) } - result = filterMetricsByVolumeState([]*workertypes.VolumeHealthMetrics{}, volumeStateFull) + result = pluginworker.FilterMetricsByVolumeState([]*workertypes.VolumeHealthMetrics{}, pluginworker.VolumeStateFull) if len(result) != 0 { t.Fatalf("expected 0 metrics for empty input, got %d", len(result)) } @@ -758,3 +758,37 @@ func workerConfigFormHasField(form *plugin_pb.ConfigForm, fieldName string) bool } return false } + +type noopDetectionSender struct{} + +func (noopDetectionSender) SendProposals(*plugin_pb.DetectionProposals) error { return nil } +func (noopDetectionSender) SendComplete(*plugin_pb.DetectionComplete) error { return nil } +func (noopDetectionSender) SendActivity(*plugin_pb.ActivityEvent) error { return nil } + +type noopExecutionSender struct{} + +func (noopExecutionSender) SendProgress(*plugin_pb.JobProgressUpdate) error { return nil } +func (noopExecutionSender) SendCompleted(*plugin_pb.JobCompleted) error { return nil } + +type recordingDetectionSender struct { + proposals *plugin_pb.DetectionProposals + complete *plugin_pb.DetectionComplete + events []*plugin_pb.ActivityEvent +} + +func (r *recordingDetectionSender) SendProposals(proposals *plugin_pb.DetectionProposals) error { + r.proposals = proposals + return nil +} + +func (r *recordingDetectionSender) SendComplete(complete *plugin_pb.DetectionComplete) error { + r.complete = complete + return nil +} + +func (r *recordingDetectionSender) SendActivity(event *plugin_pb.ActivityEvent) error { + if event != nil { + r.events = append(r.events, event) + } + return nil +} diff --git a/weed/plugin/worker/ec_balance_handler.go b/weed/worker/tasks/ec_balance/plugin_handler.go similarity index 88% rename from weed/plugin/worker/ec_balance_handler.go rename to weed/worker/tasks/ec_balance/plugin_handler.go index 092bf5ac4..5015c5bce 100644 --- a/weed/plugin/worker/ec_balance_handler.go +++ b/weed/worker/tasks/ec_balance/plugin_handler.go @@ -1,4 +1,4 @@ -package pluginworker +package ec_balance import ( "context" @@ -11,8 +11,8 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" "github.com/seaweedfs/seaweedfs/weed/util" - ecbalancetask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/ec_balance" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/grpc" "google.golang.org/protobuf/proto" @@ -25,18 +25,18 @@ const ( ) func init() { - RegisterHandler(HandlerFactory{ + pluginworker.RegisterHandler(pluginworker.HandlerFactory{ JobType: "ec_balance", - Category: CategoryDefault, + Category: pluginworker.CategoryDefault, Aliases: []string{"ec-balance", "ec.balance", "ec_shard_balance"}, - Build: func(opts HandlerBuildOptions) (JobHandler, error) { + Build: func(opts pluginworker.HandlerBuildOptions) (pluginworker.JobHandler, error) { return NewECBalanceHandler(opts.GrpcDialOption), nil }, }) } type ecBalanceWorkerConfig struct { - TaskConfig *ecbalancetask.Config + TaskConfig *Config } // ECBalanceHandler is the plugin job handler for EC shard balancing. @@ -180,7 +180,7 @@ func (h *ECBalanceHandler) Descriptor() *plugin_pb.JobTypeDescriptor { func (h *ECBalanceHandler) Detect( ctx context.Context, request *plugin_pb.RunDetectionRequest, - sender DetectionSender, + sender pluginworker.DetectionSender, ) error { if request == nil { return fmt.Errorf("run detection request is nil") @@ -195,15 +195,15 @@ func (h *ECBalanceHandler) Detect( workerConfig := deriveECBalanceWorkerConfig(request.GetWorkerConfigValues()) // Apply admin-side scope filters - collectionFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) + collectionFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) if collectionFilter != "" { workerConfig.TaskConfig.CollectionFilter = collectionFilter } - dcFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) + dcFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) if dcFilter != "" { workerConfig.TaskConfig.DataCenterFilter = dcFilter } - diskType := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "disk_type", "")) + diskType := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "disk_type", "")) if diskType != "" { workerConfig.TaskConfig.DiskType = diskType } @@ -224,7 +224,7 @@ func (h *ECBalanceHandler) Detect( maxResults = 0 } - results, hasMore, err := ecbalancetask.Detection(ctx, metrics, clusterInfo, workerConfig.TaskConfig, maxResults) + results, hasMore, err := Detection(ctx, metrics, clusterInfo, workerConfig.TaskConfig, maxResults) if err != nil { return err } @@ -261,7 +261,7 @@ func (h *ECBalanceHandler) Detect( func (h *ECBalanceHandler) Execute( ctx context.Context, request *plugin_pb.ExecuteJobRequest, - sender ExecutionSender, + sender pluginworker.ExecutionSender, ) error { if request == nil || request.Job == nil { return fmt.Errorf("execute request/job is nil") @@ -285,7 +285,7 @@ func (h *ECBalanceHandler) Execute( return fmt.Errorf("ec balance target node is required") } - task := ecbalancetask.NewECBalanceTask( + task := NewECBalanceTask( request.Job.JobId, params.VolumeId, params.Collection, @@ -308,7 +308,7 @@ func (h *ECBalanceHandler) Execute( Stage: stage, Message: message, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity(stage, message), + pluginworker.BuildExecutorActivity(stage, message), }, }); err != nil { glog.Warningf("EC balance job %s (%s): failed to send progress (%.0f%%, stage=%q): %v, cancelling execution", @@ -325,7 +325,7 @@ func (h *ECBalanceHandler) Execute( Stage: "assigned", Message: "ec balance job accepted", Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("assigned", "ec balance job accepted"), + pluginworker.BuildExecutorActivity("assigned", "ec balance job accepted"), }, }); err != nil { return err @@ -340,7 +340,7 @@ func (h *ECBalanceHandler) Execute( Stage: "failed", Message: err.Error(), Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("failed", err.Error()), + pluginworker.BuildExecutorActivity("failed", err.Error()), }, }) return err @@ -370,7 +370,7 @@ func (h *ECBalanceHandler) Execute( }, }, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("completed", resultSummary), + pluginworker.BuildExecutorActivity("completed", resultSummary), }, }) } @@ -380,14 +380,14 @@ func (h *ECBalanceHandler) collectVolumeMetrics( masterAddresses []string, collectionFilter string, ) ([]*workertypes.VolumeHealthMetrics, *topology.ActiveTopology, error) { - metrics, activeTopology, _, err := collectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) + metrics, activeTopology, _, err := pluginworker.CollectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) return metrics, activeTopology, err } func deriveECBalanceWorkerConfig(values map[string]*plugin_pb.ConfigValue) *ecBalanceWorkerConfig { - taskConfig := ecbalancetask.NewDefaultConfig() + taskConfig := NewDefaultConfig() - imbalanceThreshold := readDoubleConfig(values, "imbalance_threshold", taskConfig.ImbalanceThreshold) + imbalanceThreshold := pluginworker.ReadDoubleConfig(values, "imbalance_threshold", taskConfig.ImbalanceThreshold) if imbalanceThreshold < ecBalanceMinImbalanceThreshold { imbalanceThreshold = ecBalanceMinImbalanceThreshold } @@ -396,16 +396,13 @@ func deriveECBalanceWorkerConfig(values map[string]*plugin_pb.ConfigValue) *ecBa } taskConfig.ImbalanceThreshold = imbalanceThreshold - minServerCountRaw := readInt64Config(values, "min_server_count", int64(taskConfig.MinServerCount)) - if minServerCountRaw < int64(ecBalanceMinServerCount) { - minServerCountRaw = int64(ecBalanceMinServerCount) + minServerCount := pluginworker.ReadIntConfig(values, "min_server_count", taskConfig.MinServerCount) + if minServerCount < ecBalanceMinServerCount { + minServerCount = ecBalanceMinServerCount } - if minServerCountRaw > math.MaxInt32 { - minServerCountRaw = math.MaxInt32 - } - taskConfig.MinServerCount = int(minServerCountRaw) + taskConfig.MinServerCount = minServerCount - taskConfig.PreferredTags = util.NormalizeTagList(readStringListConfig(values, "preferred_tags")) + taskConfig.PreferredTags = util.NormalizeTagList(pluginworker.ReadStringListConfig(values, "preferred_tags")) return &ecBalanceWorkerConfig{ TaskConfig: taskConfig, @@ -467,7 +464,7 @@ func buildECBalanceProposal(result *workertypes.TaskDetectionResult) (*plugin_pb ProposalId: proposalID, DedupeKey: dedupeKey, JobType: "ec_balance", - Priority: mapTaskPriority(result.Priority), + Priority: pluginworker.MapTaskPriority(result.Priority), Summary: summary, Detail: strings.TrimSpace(result.Reason), Parameters: map[string]*plugin_pb.ConfigValue{ @@ -506,7 +503,7 @@ func decodeECBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, e } // Try protobuf-encoded params first (preferred path) - if payload := readBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { + if payload := pluginworker.ReadBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { params := &worker_pb.TaskParams{} if err := proto.Unmarshal(payload, params); err != nil { return nil, fmt.Errorf("decodeECBalanceTaskParams: unmarshal task_params_pb: %w", err) @@ -532,10 +529,10 @@ func decodeECBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, e // Legacy fallback: construct TaskParams from individual scalar parameters. // All execution-critical fields are required. - volumeID := readInt64Config(job.Parameters, "volume_id", 0) - sourceNode := strings.TrimSpace(readStringConfig(job.Parameters, "source_server", "")) - targetNode := strings.TrimSpace(readStringConfig(job.Parameters, "target_server", "")) - collection := readStringConfig(job.Parameters, "collection", "") + volumeID := pluginworker.ReadInt64Config(job.Parameters, "volume_id", 0) + sourceNode := strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "source_server", "")) + targetNode := strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "target_server", "")) + collection := pluginworker.ReadStringConfig(job.Parameters, "collection", "") if volumeID <= 0 || volumeID > math.MaxUint32 { return nil, fmt.Errorf("decodeECBalanceTaskParams: invalid or missing volume_id: %d", volumeID) @@ -551,16 +548,16 @@ func decodeECBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, e if !hasShardID || shardIDVal == nil { return nil, fmt.Errorf("decodeECBalanceTaskParams: missing shard_id (required for EcBalanceTaskParams)") } - shardID := readInt64Config(job.Parameters, "shard_id", -1) + shardID := pluginworker.ReadInt64Config(job.Parameters, "shard_id", -1) if shardID < 0 || shardID > math.MaxUint32 { return nil, fmt.Errorf("decodeECBalanceTaskParams: invalid shard_id: %d", shardID) } - sourceDiskID := readInt64Config(job.Parameters, "source_disk_id", 0) + sourceDiskID := pluginworker.ReadInt64Config(job.Parameters, "source_disk_id", 0) if sourceDiskID < 0 || sourceDiskID > math.MaxUint32 { return nil, fmt.Errorf("decodeECBalanceTaskParams: invalid source_disk_id: %d", sourceDiskID) } - targetDiskID := readInt64Config(job.Parameters, "target_disk_id", 0) + targetDiskID := pluginworker.ReadInt64Config(job.Parameters, "target_disk_id", 0) if targetDiskID < 0 || targetDiskID > math.MaxUint32 { return nil, fmt.Errorf("decodeECBalanceTaskParams: invalid target_disk_id: %d", targetDiskID) } @@ -588,8 +585,8 @@ func decodeECBalanceTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, e } func emitECBalanceDecisionTrace( - sender DetectionSender, - taskConfig *ecbalancetask.Config, + sender pluginworker.DetectionSender, + taskConfig *Config, results []*workertypes.TaskDetectionResult, maxResults int, hasMore bool, @@ -627,7 +624,7 @@ func emitECBalanceDecisionTrace( phaseCounts["global"], ) - return sender.SendActivity(BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ + return sender.SendActivity(pluginworker.BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ "total_moves": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(len(results))}, }, diff --git a/weed/plugin/worker/ec_balance_handler_test.go b/weed/worker/tasks/ec_balance/plugin_handler_test.go similarity index 98% rename from weed/plugin/worker/ec_balance_handler_test.go rename to weed/worker/tasks/ec_balance/plugin_handler_test.go index d5583e50b..5c3b8c96e 100644 --- a/weed/plugin/worker/ec_balance_handler_test.go +++ b/weed/worker/tasks/ec_balance/plugin_handler_test.go @@ -1,11 +1,10 @@ -package pluginworker +package ec_balance import ( "testing" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" - ecbalancetask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/ec_balance" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/protobuf/proto" ) @@ -290,7 +289,7 @@ func TestECBalanceHandlerCapability(t *testing.T) { } func TestECBalanceConfigRoundTrip(t *testing.T) { - config := ecbalancetask.NewDefaultConfig() + config := NewDefaultConfig() config.ImbalanceThreshold = 0.3 config.MinServerCount = 5 config.CollectionFilter = "my_col" @@ -302,7 +301,7 @@ func TestECBalanceConfigRoundTrip(t *testing.T) { t.Fatal("expected non-nil policy") } - config2 := ecbalancetask.NewDefaultConfig() + config2 := NewDefaultConfig() if err := config2.FromTaskPolicy(policy); err != nil { t.Fatalf("failed to load from policy: %v", err) } diff --git a/weed/plugin/worker/erasure_coding_handler.go b/weed/worker/tasks/erasure_coding/plugin_handler.go similarity index 85% rename from weed/plugin/worker/erasure_coding_handler.go rename to weed/worker/tasks/erasure_coding/plugin_handler.go index abd7fd5ca..bc659828a 100644 --- a/weed/plugin/worker/erasure_coding_handler.go +++ b/weed/worker/tasks/erasure_coding/plugin_handler.go @@ -1,4 +1,4 @@ -package pluginworker +package erasure_coding import ( "context" @@ -11,28 +11,28 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" ecstorage "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" "github.com/seaweedfs/seaweedfs/weed/util" "github.com/seaweedfs/seaweedfs/weed/util/wildcard" - erasurecodingtask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/grpc" "google.golang.org/protobuf/proto" ) func init() { - RegisterHandler(HandlerFactory{ + pluginworker.RegisterHandler(pluginworker.HandlerFactory{ JobType: "erasure_coding", - Category: CategoryHeavy, + Category: pluginworker.CategoryHeavy, Aliases: []string{"erasure-coding", "erasure.coding", "ec"}, - Build: func(opts HandlerBuildOptions) (JobHandler, error) { + Build: func(opts pluginworker.HandlerBuildOptions) (pluginworker.JobHandler, error) { return NewErasureCodingHandler(opts.GrpcDialOption, opts.WorkingDir), nil }, }) } type erasureCodingWorkerConfig struct { - TaskConfig *erasurecodingtask.Config + TaskConfig *Config } // ErasureCodingHandler is the plugin job handler for erasure coding. @@ -188,7 +188,7 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor { func (h *ErasureCodingHandler) Detect( ctx context.Context, request *plugin_pb.RunDetectionRequest, - sender DetectionSender, + sender pluginworker.DetectionSender, ) error { if request == nil { return fmt.Errorf("run detection request is nil") @@ -202,7 +202,7 @@ func (h *ErasureCodingHandler) Detect( workerConfig := deriveErasureCodingWorkerConfig(request.GetWorkerConfigValues()) - collectionFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) + collectionFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) if collectionFilter != "" { workerConfig.TaskConfig.CollectionFilter = collectionFilter } @@ -222,7 +222,7 @@ func (h *ErasureCodingHandler) Detect( if maxResults < 0 { maxResults = 0 } - results, hasMore, err := erasurecodingtask.Detection(ctx, metrics, clusterInfo, workerConfig.TaskConfig, maxResults) + results, hasMore, err := Detection(ctx, metrics, clusterInfo, workerConfig.TaskConfig, maxResults) if err != nil { return err } @@ -256,9 +256,9 @@ func (h *ErasureCodingHandler) Detect( } func emitErasureCodingDetectionDecisionTrace( - sender DetectionSender, + sender pluginworker.DetectionSender, metrics []*workertypes.VolumeHealthMetrics, - taskConfig *erasurecodingtask.Config, + taskConfig *Config, results []*workertypes.TaskDetectionResult, maxResults int, hasMore bool, @@ -352,7 +352,7 @@ func emitErasureCodingDetectionDecisionTrace( ) } - if err := sender.SendActivity(BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_summary", summaryMessage, map[string]*plugin_pb.ConfigValue{ "total_volumes": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(totalVolumes)}, }, @@ -409,7 +409,7 @@ func emitErasureCodingDetectionDecisionTrace( metric.FullnessRatio*100, taskConfig.FullnessRatio*100, ) - if err := sender.SendActivity(BuildDetectorActivity("decision_volume", message, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_volume", message, map[string]*plugin_pb.ConfigValue{ "volume_id": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(metric.VolumeID)}, }, @@ -446,7 +446,7 @@ func emitErasureCodingDetectionDecisionTrace( func (h *ErasureCodingHandler) Execute( ctx context.Context, request *plugin_pb.ExecuteJobRequest, - sender ExecutionSender, + sender pluginworker.ExecutionSender, ) error { if request == nil || request.Job == nil { return fmt.Errorf("execute request/job is nil") @@ -472,7 +472,7 @@ func (h *ErasureCodingHandler) Execute( return fmt.Errorf("erasure coding targets are required") } - task := erasurecodingtask.NewErasureCodingTask( + task := NewErasureCodingTask( request.Job.JobId, params.Sources[0].Node, params.VolumeId, @@ -494,7 +494,7 @@ func (h *ErasureCodingHandler) Execute( Stage: stage, Message: message, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity(stage, message), + pluginworker.BuildExecutorActivity(stage, message), }, }); err != nil { execCancel() @@ -509,7 +509,7 @@ func (h *ErasureCodingHandler) Execute( Stage: "assigned", Message: "erasure coding job accepted", Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("assigned", "erasure coding job accepted"), + pluginworker.BuildExecutorActivity("assigned", "erasure coding job accepted"), }, }); err != nil { return err @@ -524,7 +524,7 @@ func (h *ErasureCodingHandler) Execute( Stage: "failed", Message: err.Error(), Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("failed", err.Error()), + pluginworker.BuildExecutorActivity("failed", err.Error()), }, }) return err @@ -552,7 +552,7 @@ func (h *ErasureCodingHandler) Execute( }, }, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("completed", resultSummary), + pluginworker.BuildExecutorActivity("completed", resultSummary), }, }) } @@ -562,20 +562,20 @@ func (h *ErasureCodingHandler) collectVolumeMetrics( masterAddresses []string, collectionFilter string, ) ([]*workertypes.VolumeHealthMetrics, *topology.ActiveTopology, error) { - metrics, activeTopology, _, err := collectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) + metrics, activeTopology, _, err := pluginworker.CollectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) return metrics, activeTopology, err } func deriveErasureCodingWorkerConfig(values map[string]*plugin_pb.ConfigValue) *erasureCodingWorkerConfig { - taskConfig := erasurecodingtask.NewDefaultConfig() + taskConfig := NewDefaultConfig() - quietForSeconds := int(readInt64Config(values, "quiet_for_seconds", int64(taskConfig.QuietForSeconds))) + quietForSeconds := pluginworker.ReadIntConfig(values, "quiet_for_seconds", taskConfig.QuietForSeconds) if quietForSeconds < 0 { quietForSeconds = 0 } taskConfig.QuietForSeconds = quietForSeconds - fullnessRatio := readDoubleConfig(values, "fullness_ratio", taskConfig.FullnessRatio) + fullnessRatio := pluginworker.ReadDoubleConfig(values, "fullness_ratio", taskConfig.FullnessRatio) if fullnessRatio < 0 { fullnessRatio = 0 } @@ -584,13 +584,13 @@ func deriveErasureCodingWorkerConfig(values map[string]*plugin_pb.ConfigValue) * } taskConfig.FullnessRatio = fullnessRatio - minSizeMB := int(readInt64Config(values, "min_size_mb", int64(taskConfig.MinSizeMB))) + minSizeMB := pluginworker.ReadIntConfig(values, "min_size_mb", taskConfig.MinSizeMB) if minSizeMB < 1 { minSizeMB = 1 } taskConfig.MinSizeMB = minSizeMB - taskConfig.PreferredTags = util.NormalizeTagList(readStringListConfig(values, "preferred_tags")) + taskConfig.PreferredTags = util.NormalizeTagList(pluginworker.ReadStringListConfig(values, "preferred_tags")) return &erasureCodingWorkerConfig{ TaskConfig: taskConfig, @@ -645,7 +645,7 @@ func buildErasureCodingProposal( ProposalId: proposalID, DedupeKey: dedupeKey, JobType: "erasure_coding", - Priority: mapTaskPriority(result.Priority), + Priority: pluginworker.MapTaskPriority(result.Priority), Summary: summary, Detail: strings.TrimSpace(result.Reason), Parameters: map[string]*plugin_pb.ConfigValue{ @@ -683,7 +683,7 @@ func decodeErasureCodingTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return nil, fmt.Errorf("job spec is nil") } - if payload := readBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { + if payload := pluginworker.ReadBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { params := &worker_pb.TaskParams{} if err := proto.Unmarshal(payload, params); err != nil { return nil, fmt.Errorf("unmarshal task_params_pb: %w", err) @@ -694,28 +694,28 @@ func decodeErasureCodingTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return params, nil } - volumeID := readInt64Config(job.Parameters, "volume_id", 0) - sourceNode := strings.TrimSpace(readStringConfig(job.Parameters, "source_server", "")) + volumeID := pluginworker.ReadUint32Config(job.Parameters, "volume_id", 0) + sourceNode := strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "source_server", "")) if sourceNode == "" { - sourceNode = strings.TrimSpace(readStringConfig(job.Parameters, "server", "")) + sourceNode = strings.TrimSpace(pluginworker.ReadStringConfig(job.Parameters, "server", "")) } - targetServers := readStringListConfig(job.Parameters, "target_servers") + targetServers := pluginworker.ReadStringListConfig(job.Parameters, "target_servers") if len(targetServers) == 0 { - targetServers = readStringListConfig(job.Parameters, "targets") + targetServers = pluginworker.ReadStringListConfig(job.Parameters, "targets") } - collection := readStringConfig(job.Parameters, "collection", "") + collection := pluginworker.ReadStringConfig(job.Parameters, "collection", "") - dataShards := int32(readInt64Config(job.Parameters, "data_shards", int64(ecstorage.DataShardsCount))) + dataShards := pluginworker.ReadInt32Config(job.Parameters, "data_shards", int32(ecstorage.DataShardsCount)) if dataShards <= 0 { - dataShards = ecstorage.DataShardsCount + dataShards = int32(ecstorage.DataShardsCount) } - parityShards := int32(readInt64Config(job.Parameters, "parity_shards", int64(ecstorage.ParityShardsCount))) + parityShards := pluginworker.ReadInt32Config(job.Parameters, "parity_shards", int32(ecstorage.ParityShardsCount)) if parityShards <= 0 { - parityShards = ecstorage.ParityShardsCount + parityShards = int32(ecstorage.ParityShardsCount) } totalShards := int(dataShards + parityShards) - if volumeID <= 0 { + if volumeID == 0 { return nil, fmt.Errorf("missing volume_id in job parameters") } if sourceNode == "" { @@ -737,7 +737,7 @@ func decodeErasureCodingTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam } targets = append(targets, &worker_pb.TaskTarget{ Node: targetNode, - VolumeId: uint32(volumeID), + VolumeId: volumeID, ShardIds: shardAssignments[i], }) } @@ -747,12 +747,12 @@ func decodeErasureCodingTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParam return &worker_pb.TaskParams{ TaskId: job.JobId, - VolumeId: uint32(volumeID), + VolumeId: volumeID, Collection: collection, Sources: []*worker_pb.TaskSource{ { Node: sourceNode, - VolumeId: uint32(volumeID), + VolumeId: volumeID, }, }, Targets: targets, @@ -819,71 +819,6 @@ func applyErasureCodingExecutionDefaults( } } -func readStringListConfig(values map[string]*plugin_pb.ConfigValue, field string) []string { - if values == nil { - return nil - } - value := values[field] - if value == nil { - return nil - } - - switch kind := value.Kind.(type) { - case *plugin_pb.ConfigValue_StringList: - return normalizeStringList(kind.StringList.GetValues()) - case *plugin_pb.ConfigValue_ListValue: - out := make([]string, 0, len(kind.ListValue.GetValues())) - for _, item := range kind.ListValue.GetValues() { - itemText := readStringFromConfigValue(item) - if itemText != "" { - out = append(out, itemText) - } - } - return normalizeStringList(out) - case *plugin_pb.ConfigValue_StringValue: - return normalizeStringList(strings.Split(kind.StringValue, ",")) - } - - return nil -} - -func readStringFromConfigValue(value *plugin_pb.ConfigValue) string { - if value == nil { - return "" - } - switch kind := value.Kind.(type) { - case *plugin_pb.ConfigValue_StringValue: - return strings.TrimSpace(kind.StringValue) - case *plugin_pb.ConfigValue_Int64Value: - return fmt.Sprintf("%d", kind.Int64Value) - case *plugin_pb.ConfigValue_DoubleValue: - return fmt.Sprintf("%g", kind.DoubleValue) - case *plugin_pb.ConfigValue_BoolValue: - if kind.BoolValue { - return "true" - } - return "false" - } - return "" -} - -func normalizeStringList(values []string) []string { - normalized := make([]string, 0, len(values)) - seen := make(map[string]struct{}, len(values)) - for _, value := range values { - item := strings.TrimSpace(value) - if item == "" { - continue - } - if _, found := seen[item]; found { - continue - } - seen[item] = struct{}{} - normalized = append(normalized, item) - } - return normalized -} - func assignECShardIDs(totalShards int, targetCount int) [][]uint32 { if targetCount <= 0 { return nil diff --git a/weed/plugin/worker/erasure_coding_handler_test.go b/weed/worker/tasks/erasure_coding/plugin_handler_test.go similarity index 85% rename from weed/plugin/worker/erasure_coding_handler_test.go rename to weed/worker/tasks/erasure_coding/plugin_handler_test.go index 39b3b9fe8..6febba84b 100644 --- a/weed/plugin/worker/erasure_coding_handler_test.go +++ b/weed/worker/tasks/erasure_coding/plugin_handler_test.go @@ -1,4 +1,4 @@ -package pluginworker +package erasure_coding import ( "context" @@ -9,7 +9,6 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" ecstorage "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" - erasurecodingtask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/erasure_coding" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/protobuf/proto" ) @@ -206,7 +205,7 @@ func TestErasureCodingHandlerRejectsUnsupportedJobType(t *testing.T) { func TestEmitErasureCodingDetectionDecisionTraceNoTasks(t *testing.T) { sender := &recordingDetectionSender{} - config := erasurecodingtask.NewDefaultConfig() + config := NewDefaultConfig() config.QuietForSeconds = 5 * 60 config.MinSizeMB = 30 config.FullnessRatio = 0.91 @@ -291,3 +290,54 @@ func TestApplyErasureCodingExecutionDefaultsForcesLocalFields(t *testing.T) { t.Fatalf("expected cleanup_source true") } } + +type noopDetectionSender struct{} + +func (noopDetectionSender) SendProposals(*plugin_pb.DetectionProposals) error { return nil } +func (noopDetectionSender) SendComplete(*plugin_pb.DetectionComplete) error { return nil } +func (noopDetectionSender) SendActivity(*plugin_pb.ActivityEvent) error { return nil } + +type noopExecutionSender struct{} + +func (noopExecutionSender) SendProgress(*plugin_pb.JobProgressUpdate) error { return nil } +func (noopExecutionSender) SendCompleted(*plugin_pb.JobCompleted) error { return nil } + +type recordingDetectionSender struct { + proposals *plugin_pb.DetectionProposals + complete *plugin_pb.DetectionComplete + events []*plugin_pb.ActivityEvent +} + +func (r *recordingDetectionSender) SendProposals(proposals *plugin_pb.DetectionProposals) error { + r.proposals = proposals + return nil +} + +func (r *recordingDetectionSender) SendComplete(complete *plugin_pb.DetectionComplete) error { + r.complete = complete + return nil +} + +func (r *recordingDetectionSender) SendActivity(event *plugin_pb.ActivityEvent) error { + if event != nil { + r.events = append(r.events, event) + } + return nil +} + +func workerConfigFormHasField(form *plugin_pb.ConfigForm, fieldName string) bool { + if form == nil { + return false + } + for _, section := range form.Sections { + if section == nil { + continue + } + for _, field := range section.Fields { + if field != nil && field.Name == fieldName { + return true + } + } + } + return false +} diff --git a/weed/plugin/worker/iceberg/compact.go b/weed/worker/tasks/iceberg/compact.go similarity index 100% rename from weed/plugin/worker/iceberg/compact.go rename to weed/worker/tasks/iceberg/compact.go diff --git a/weed/plugin/worker/iceberg/config.go b/weed/worker/tasks/iceberg/config.go similarity index 100% rename from weed/plugin/worker/iceberg/config.go rename to weed/worker/tasks/iceberg/config.go diff --git a/weed/plugin/worker/iceberg/delete_rewrite.go b/weed/worker/tasks/iceberg/delete_rewrite.go similarity index 100% rename from weed/plugin/worker/iceberg/delete_rewrite.go rename to weed/worker/tasks/iceberg/delete_rewrite.go diff --git a/weed/plugin/worker/iceberg/detection.go b/weed/worker/tasks/iceberg/detection.go similarity index 100% rename from weed/plugin/worker/iceberg/detection.go rename to weed/worker/tasks/iceberg/detection.go diff --git a/weed/plugin/worker/iceberg/exec_test.go b/weed/worker/tasks/iceberg/exec_test.go similarity index 100% rename from weed/plugin/worker/iceberg/exec_test.go rename to weed/worker/tasks/iceberg/exec_test.go diff --git a/weed/plugin/worker/iceberg/filer_io.go b/weed/worker/tasks/iceberg/filer_io.go similarity index 100% rename from weed/plugin/worker/iceberg/filer_io.go rename to weed/worker/tasks/iceberg/filer_io.go diff --git a/weed/plugin/worker/iceberg/handler.go b/weed/worker/tasks/iceberg/handler.go similarity index 100% rename from weed/plugin/worker/iceberg/handler.go rename to weed/worker/tasks/iceberg/handler.go diff --git a/weed/plugin/worker/iceberg/handler_test.go b/weed/worker/tasks/iceberg/handler_test.go similarity index 100% rename from weed/plugin/worker/iceberg/handler_test.go rename to weed/worker/tasks/iceberg/handler_test.go diff --git a/weed/plugin/worker/iceberg/operations.go b/weed/worker/tasks/iceberg/operations.go similarity index 100% rename from weed/plugin/worker/iceberg/operations.go rename to weed/worker/tasks/iceberg/operations.go diff --git a/weed/plugin/worker/iceberg/planning_index.go b/weed/worker/tasks/iceberg/planning_index.go similarity index 100% rename from weed/plugin/worker/iceberg/planning_index.go rename to weed/worker/tasks/iceberg/planning_index.go diff --git a/weed/plugin/worker/iceberg/platform_test.go b/weed/worker/tasks/iceberg/platform_test.go similarity index 100% rename from weed/plugin/worker/iceberg/platform_test.go rename to weed/worker/tasks/iceberg/platform_test.go diff --git a/weed/plugin/worker/iceberg/resource_groups.go b/weed/worker/tasks/iceberg/resource_groups.go similarity index 100% rename from weed/plugin/worker/iceberg/resource_groups.go rename to weed/worker/tasks/iceberg/resource_groups.go diff --git a/weed/plugin/worker/iceberg/sort_strategy.go b/weed/worker/tasks/iceberg/sort_strategy.go similarity index 100% rename from weed/plugin/worker/iceberg/sort_strategy.go rename to weed/worker/tasks/iceberg/sort_strategy.go diff --git a/weed/plugin/worker/iceberg/testing_api.go b/weed/worker/tasks/iceberg/testing_api.go similarity index 100% rename from weed/plugin/worker/iceberg/testing_api.go rename to weed/worker/tasks/iceberg/testing_api.go diff --git a/weed/plugin/worker/iceberg/where_filter.go b/weed/worker/tasks/iceberg/where_filter.go similarity index 100% rename from weed/plugin/worker/iceberg/where_filter.go rename to weed/worker/tasks/iceberg/where_filter.go diff --git a/weed/plugin/worker/iceberg/where_filter_test.go b/weed/worker/tasks/iceberg/where_filter_test.go similarity index 100% rename from weed/plugin/worker/iceberg/where_filter_test.go rename to weed/worker/tasks/iceberg/where_filter_test.go diff --git a/weed/plugin/worker/vacuum_handler.go b/weed/worker/tasks/vacuum/plugin_handler.go similarity index 73% rename from weed/plugin/worker/vacuum_handler.go rename to weed/worker/tasks/vacuum/plugin_handler.go index 3ffa95170..1d95f61ea 100644 --- a/weed/plugin/worker/vacuum_handler.go +++ b/weed/worker/tasks/vacuum/plugin_handler.go @@ -1,9 +1,8 @@ -package pluginworker +package vacuum import ( "context" "fmt" - "strconv" "strings" "time" @@ -11,11 +10,10 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" - vacuumtask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/grpc" "google.golang.org/protobuf/proto" - "google.golang.org/protobuf/types/known/timestamppb" ) const ( @@ -24,11 +22,11 @@ const ( ) func init() { - RegisterHandler(HandlerFactory{ + pluginworker.RegisterHandler(pluginworker.HandlerFactory{ JobType: "vacuum", - Category: CategoryDefault, + Category: pluginworker.CategoryDefault, Aliases: []string{"vol.vacuum", "volume.vacuum"}, - Build: func(opts HandlerBuildOptions) (JobHandler, error) { + Build: func(opts pluginworker.HandlerBuildOptions) (pluginworker.JobHandler, error) { return NewVacuumHandler(opts.GrpcDialOption, int32(opts.MaxExecute)), nil }, }) @@ -99,9 +97,9 @@ func (h *VacuumHandler) Descriptor() *plugin_pb.JobTypeDescriptor { FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_ENUM, Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_SELECT, Options: []*plugin_pb.ConfigOption{ - {Value: string(volumeStateAll), Label: "All Volumes"}, - {Value: string(volumeStateActive), Label: "Active (writable)"}, - {Value: string(volumeStateFull), Label: "Full (read-only)"}, + {Value: string(pluginworker.VolumeStateAll), Label: "All Volumes"}, + {Value: string(pluginworker.VolumeStateActive), Label: "Active (writable)"}, + {Value: string(pluginworker.VolumeStateFull), Label: "Full (read-only)"}, }, }, { @@ -136,7 +134,7 @@ func (h *VacuumHandler) Descriptor() *plugin_pb.JobTypeDescriptor { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""}, }, "volume_state": { - Kind: &plugin_pb.ConfigValue_StringValue{StringValue: string(volumeStateAll)}, + Kind: &plugin_pb.ConfigValue_StringValue{StringValue: string(pluginworker.VolumeStateAll)}, }, "data_center_filter": { Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""}, @@ -198,7 +196,7 @@ func (h *VacuumHandler) Descriptor() *plugin_pb.JobTypeDescriptor { } } -func (h *VacuumHandler) Detect(ctx context.Context, request *plugin_pb.RunDetectionRequest, sender DetectionSender) error { +func (h *VacuumHandler) Detect(ctx context.Context, request *plugin_pb.RunDetectionRequest, sender pluginworker.DetectionSender) error { if request == nil { return fmt.Errorf("run detection request is nil") } @@ -210,7 +208,7 @@ func (h *VacuumHandler) Detect(ctx context.Context, request *plugin_pb.RunDetect } workerConfig := deriveVacuumConfig(request.GetWorkerConfigValues()) - collectionFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) + collectionFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "collection_filter", "")) masters := make([]string, 0) if request.ClusterContext != nil { masters = append(masters, request.ClusterContext.MasterGrpcAddresses...) @@ -220,19 +218,19 @@ func (h *VacuumHandler) Detect(ctx context.Context, request *plugin_pb.RunDetect return err } - volState := volumeState(strings.ToUpper(strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "volume_state", string(volumeStateAll))))) - metrics = filterMetricsByVolumeState(metrics, volState) + volState := pluginworker.VolumeState(strings.ToUpper(strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "volume_state", string(pluginworker.VolumeStateAll))))) + metrics = pluginworker.FilterMetricsByVolumeState(metrics, volState) - dataCenterFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) - rackFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "rack_filter", "")) - nodeFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "node_filter", "")) + dataCenterFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "data_center_filter", "")) + rackFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "rack_filter", "")) + nodeFilter := strings.TrimSpace(pluginworker.ReadStringConfig(request.GetAdminConfigValues(), "node_filter", "")) if dataCenterFilter != "" || rackFilter != "" || nodeFilter != "" { - metrics = filterMetricsByLocation(metrics, dataCenterFilter, rackFilter, nodeFilter) + metrics = pluginworker.FilterMetricsByLocation(metrics, dataCenterFilter, rackFilter, nodeFilter) } clusterInfo := &workertypes.ClusterInfo{ActiveTopology: activeTopology} - results, err := vacuumtask.Detection(metrics, clusterInfo, workerConfig) + results, err := Detection(metrics, clusterInfo, workerConfig) if err != nil { return err } @@ -273,9 +271,9 @@ func (h *VacuumHandler) Detect(ctx context.Context, request *plugin_pb.RunDetect } func emitVacuumDetectionDecisionTrace( - sender DetectionSender, + sender pluginworker.DetectionSender, metrics []*workertypes.VolumeHealthMetrics, - workerConfig *vacuumtask.Config, + workerConfig *Config, results []*workertypes.TaskDetectionResult, ) error { if sender == nil || workerConfig == nil { @@ -313,7 +311,7 @@ func emitVacuumDetectionDecisionTrace( ) } - if err := sender.SendActivity(BuildDetectorActivity(summaryStage, summaryMessage, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity(summaryStage, summaryMessage, map[string]*plugin_pb.ConfigValue{ "total_volumes": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(totalVolumes)}, }, @@ -345,7 +343,7 @@ func emitVacuumDetectionDecisionTrace( metric.GarbageRatio*100, workerConfig.GarbageThreshold*100, ) - if err := sender.SendActivity(BuildDetectorActivity("decision_volume", message, map[string]*plugin_pb.ConfigValue{ + if err := sender.SendActivity(pluginworker.BuildDetectorActivity("decision_volume", message, map[string]*plugin_pb.ConfigValue{ "volume_id": { Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(metric.VolumeID)}, }, @@ -363,7 +361,7 @@ func emitVacuumDetectionDecisionTrace( return nil } -func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJobRequest, sender ExecutionSender) error { +func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJobRequest, sender pluginworker.ExecutionSender) error { if request == nil || request.Job == nil { return fmt.Errorf("execute request/job is nil") } @@ -397,7 +395,7 @@ func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJ } } - task := vacuumtask.NewVacuumTask( + task := NewVacuumTask( request.Job.JobId, params.Sources[0].Node, params.VolumeId, @@ -419,7 +417,7 @@ func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJ Stage: stage, Message: message, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity(stage, message), + pluginworker.BuildExecutorActivity(stage, message), }, }); err != nil { execCancel() @@ -434,7 +432,7 @@ func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJ Stage: "assigned", Message: "vacuum job accepted", Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("assigned", "vacuum job accepted"), + pluginworker.BuildExecutorActivity("assigned", "vacuum job accepted"), }, }); err != nil { return err @@ -449,7 +447,7 @@ func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJ Stage: "failed", Message: err.Error(), Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("failed", err.Error()), + pluginworker.BuildExecutorActivity("failed", err.Error()), }, }) return err @@ -472,7 +470,7 @@ func (h *VacuumHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJ }, }, Activities: []*plugin_pb.ActivityEvent{ - BuildExecutorActivity("completed", resultSummary), + pluginworker.BuildExecutorActivity("completed", resultSummary), }, }) } @@ -482,13 +480,13 @@ func (h *VacuumHandler) collectVolumeMetrics( masterAddresses []string, collectionFilter string, ) ([]*workertypes.VolumeHealthMetrics, *topology.ActiveTopology, error) { - metrics, activeTopology, _, err := collectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) + metrics, activeTopology, _, err := pluginworker.CollectVolumeMetricsFromMasters(ctx, masterAddresses, collectionFilter, h.grpcDialOption) return metrics, activeTopology, err } -func deriveVacuumConfig(values map[string]*plugin_pb.ConfigValue) *vacuumtask.Config { - config := vacuumtask.NewDefaultConfig() - config.GarbageThreshold = readDoubleConfig(values, "garbage_threshold", config.GarbageThreshold) +func deriveVacuumConfig(values map[string]*plugin_pb.ConfigValue) *Config { + config := NewDefaultConfig() + config.GarbageThreshold = pluginworker.ReadDoubleConfig(values, "garbage_threshold", config.GarbageThreshold) config.MinVolumeAgeSeconds = 0 // plugin worker does not filter by volume age return config } @@ -529,7 +527,7 @@ func buildVacuumProposal(result *workertypes.TaskDetectionResult) (*plugin_pb.Jo ProposalId: proposalID, DedupeKey: dedupeKey, JobType: "vacuum", - Priority: mapTaskPriority(result.Priority), + Priority: pluginworker.MapTaskPriority(result.Priority), Summary: summary, Detail: strings.TrimSpace(result.Reason), Parameters: map[string]*plugin_pb.ConfigValue{ @@ -563,7 +561,7 @@ func decodeVacuumTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, erro return nil, fmt.Errorf("job spec is nil") } - if payload := readBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { + if payload := pluginworker.ReadBytesConfig(job.Parameters, "task_params_pb"); len(payload) > 0 { params := &worker_pb.TaskParams{} if err := proto.Unmarshal(payload, params); err != nil { return nil, fmt.Errorf("unmarshal task_params_pb: %w", err) @@ -574,10 +572,10 @@ func decodeVacuumTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, erro return params, nil } - volumeID := readInt64Config(job.Parameters, "volume_id", 0) - server := readStringConfig(job.Parameters, "server", "") - collection := readStringConfig(job.Parameters, "collection", "") - if volumeID <= 0 { + volumeID := pluginworker.ReadUint32Config(job.Parameters, "volume_id", 0) + server := pluginworker.ReadStringConfig(job.Parameters, "server", "") + collection := pluginworker.ReadStringConfig(job.Parameters, "collection", "") + if volumeID == 0 { return nil, fmt.Errorf("missing volume_id in job parameters") } if strings.TrimSpace(server) == "" { @@ -586,12 +584,12 @@ func decodeVacuumTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, erro return &worker_pb.TaskParams{ TaskId: job.JobId, - VolumeId: uint32(volumeID), + VolumeId: volumeID, Collection: collection, Sources: []*worker_pb.TaskSource{ { Node: server, - VolumeId: uint32(volumeID), + VolumeId: volumeID, }, }, TaskParams: &worker_pb.TaskParams_VacuumParams{ @@ -603,141 +601,3 @@ func decodeVacuumTaskParams(job *plugin_pb.JobSpec) (*worker_pb.TaskParams, erro }, }, nil } - -func readStringConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback string) string { - if values == nil { - return fallback - } - value := values[field] - if value == nil { - return fallback - } - switch kind := value.Kind.(type) { - case *plugin_pb.ConfigValue_StringValue: - return kind.StringValue - case *plugin_pb.ConfigValue_Int64Value: - return strconv.FormatInt(kind.Int64Value, 10) - case *plugin_pb.ConfigValue_DoubleValue: - return strconv.FormatFloat(kind.DoubleValue, 'f', -1, 64) - case *plugin_pb.ConfigValue_BoolValue: - return strconv.FormatBool(kind.BoolValue) - } - return fallback -} - -func readDoubleConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback float64) float64 { - if values == nil { - return fallback - } - value := values[field] - if value == nil { - return fallback - } - switch kind := value.Kind.(type) { - case *plugin_pb.ConfigValue_DoubleValue: - return kind.DoubleValue - case *plugin_pb.ConfigValue_Int64Value: - return float64(kind.Int64Value) - case *plugin_pb.ConfigValue_StringValue: - parsed, err := strconv.ParseFloat(strings.TrimSpace(kind.StringValue), 64) - if err == nil { - return parsed - } - case *plugin_pb.ConfigValue_BoolValue: - if kind.BoolValue { - return 1 - } - return 0 - } - return fallback -} - -func readInt64Config(values map[string]*plugin_pb.ConfigValue, field string, fallback int64) int64 { - if values == nil { - return fallback - } - value := values[field] - if value == nil { - return fallback - } - switch kind := value.Kind.(type) { - case *plugin_pb.ConfigValue_Int64Value: - return kind.Int64Value - case *plugin_pb.ConfigValue_DoubleValue: - return int64(kind.DoubleValue) - case *plugin_pb.ConfigValue_StringValue: - parsed, err := strconv.ParseInt(strings.TrimSpace(kind.StringValue), 10, 64) - if err == nil { - return parsed - } - case *plugin_pb.ConfigValue_BoolValue: - if kind.BoolValue { - return 1 - } - return 0 - } - return fallback -} - -func readBytesConfig(values map[string]*plugin_pb.ConfigValue, field string) []byte { - if values == nil { - return nil - } - value := values[field] - if value == nil { - return nil - } - if kind, ok := value.Kind.(*plugin_pb.ConfigValue_BytesValue); ok { - return kind.BytesValue - } - return nil -} - -func mapTaskPriority(priority workertypes.TaskPriority) plugin_pb.JobPriority { - switch strings.ToLower(string(priority)) { - case "low": - return plugin_pb.JobPriority_JOB_PRIORITY_LOW - case "medium", "normal": - return plugin_pb.JobPriority_JOB_PRIORITY_NORMAL - case "high": - return plugin_pb.JobPriority_JOB_PRIORITY_HIGH - case "critical": - return plugin_pb.JobPriority_JOB_PRIORITY_CRITICAL - default: - return plugin_pb.JobPriority_JOB_PRIORITY_NORMAL - } -} - -// ShouldSkipDetectionByInterval returns true when less than minIntervalSeconds -// have elapsed since lastSuccessfulRun. Exported so sub-packages can reuse it. -func ShouldSkipDetectionByInterval(lastSuccessfulRun *timestamppb.Timestamp, minIntervalSeconds int) bool { - if lastSuccessfulRun == nil || minIntervalSeconds <= 0 { - return false - } - lastRun := lastSuccessfulRun.AsTime() - if lastRun.IsZero() { - return false - } - return time.Since(lastRun) < time.Duration(minIntervalSeconds)*time.Second -} - -// BuildExecutorActivity creates an executor activity event. Exported for sub-packages. -func BuildExecutorActivity(stage string, message string) *plugin_pb.ActivityEvent { - return &plugin_pb.ActivityEvent{ - Source: plugin_pb.ActivitySource_ACTIVITY_SOURCE_EXECUTOR, - Stage: stage, - Message: message, - CreatedAt: timestamppb.Now(), - } -} - -// BuildDetectorActivity creates a detector activity event. Exported for sub-packages. -func BuildDetectorActivity(stage string, message string, details map[string]*plugin_pb.ConfigValue) *plugin_pb.ActivityEvent { - return &plugin_pb.ActivityEvent{ - Source: plugin_pb.ActivitySource_ACTIVITY_SOURCE_DETECTOR, - Stage: stage, - Message: message, - Details: details, - CreatedAt: timestamppb.Now(), - } -} diff --git a/weed/plugin/worker/vacuum_handler_test.go b/weed/worker/tasks/vacuum/plugin_handler_test.go similarity index 91% rename from weed/plugin/worker/vacuum_handler_test.go rename to weed/worker/tasks/vacuum/plugin_handler_test.go index 97a0d275f..676e15efe 100644 --- a/weed/plugin/worker/vacuum_handler_test.go +++ b/weed/worker/tasks/vacuum/plugin_handler_test.go @@ -1,4 +1,4 @@ -package pluginworker +package vacuum import ( "context" @@ -8,7 +8,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb" "github.com/seaweedfs/seaweedfs/weed/pb/worker_pb" - vacuumtask "github.com/seaweedfs/seaweedfs/weed/worker/tasks/vacuum" + pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker" workertypes "github.com/seaweedfs/seaweedfs/weed/worker/types" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/timestamppb" @@ -93,7 +93,7 @@ func TestDeriveVacuumConfigAllowsZeroValues(t *testing.T) { } func TestMasterAddressCandidates(t *testing.T) { - candidates := masterAddressCandidates("localhost:9333") + candidates := pluginworker.MasterAddressCandidates("localhost:9333") if len(candidates) != 2 { t.Fatalf("expected 2 candidates, got %d: %v", len(candidates), candidates) } @@ -110,20 +110,20 @@ func TestMasterAddressCandidates(t *testing.T) { } func TestShouldSkipDetectionByInterval(t *testing.T) { - if ShouldSkipDetectionByInterval(nil, 10) { + if pluginworker.ShouldSkipDetectionByInterval(nil, 10) { t.Fatalf("expected false when timestamp is nil") } - if ShouldSkipDetectionByInterval(timestamppb.Now(), 0) { + if pluginworker.ShouldSkipDetectionByInterval(timestamppb.Now(), 0) { t.Fatalf("expected false when min interval is zero") } recent := timestamppb.New(time.Now().Add(-5 * time.Second)) - if !ShouldSkipDetectionByInterval(recent, 10) { + if !pluginworker.ShouldSkipDetectionByInterval(recent, 10) { t.Fatalf("expected true for recent successful run") } old := timestamppb.New(time.Now().Add(-30 * time.Second)) - if ShouldSkipDetectionByInterval(old, 10) { + if pluginworker.ShouldSkipDetectionByInterval(old, 10) { t.Fatalf("expected false for old successful run") } } @@ -146,7 +146,7 @@ func TestVacuumHandlerRejectsUnsupportedJobType(t *testing.T) { } func TestBuildExecutorActivity(t *testing.T) { - activity := BuildExecutorActivity("running", "vacuum in progress") + activity := pluginworker.BuildExecutorActivity("running", "vacuum in progress") if activity == nil { t.Fatalf("expected non-nil activity") } @@ -166,7 +166,7 @@ func TestBuildExecutorActivity(t *testing.T) { func TestEmitVacuumDetectionDecisionTraceNoTasks(t *testing.T) { sender := &recordingDetectionSender{} - config := vacuumtask.NewDefaultConfig() + config := NewDefaultConfig() config.GarbageThreshold = 0.3 config.MinVolumeAgeSeconds = 0 @@ -279,17 +279,17 @@ func TestVacuumFiltersVolumeState(t *testing.T) { tests := []struct { name string - state volumeState + state pluginworker.VolumeState wantIDs []uint32 }{ - {"ALL returns all", volumeStateAll, []uint32{1, 2, 3, 4}}, - {"ACTIVE returns writable", volumeStateActive, []uint32{1, 3}}, - {"FULL returns read-only", volumeStateFull, []uint32{2, 4}}, + {"ALL returns all", pluginworker.VolumeStateAll, []uint32{1, 2, 3, 4}}, + {"ACTIVE returns writable", pluginworker.VolumeStateActive, []uint32{1, 3}}, + {"FULL returns read-only", pluginworker.VolumeStateFull, []uint32{2, 4}}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - filtered := filterMetricsByVolumeState(metrics, tt.state) + filtered := pluginworker.FilterMetricsByVolumeState(metrics, tt.state) if len(filtered) != len(tt.wantIDs) { t.Fatalf("got %d metrics, want %d", len(filtered), len(tt.wantIDs)) } @@ -325,7 +325,7 @@ func TestVacuumFiltersLocation(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - filtered := filterMetricsByLocation(metrics, tt.dc, tt.rack, tt.node) + filtered := pluginworker.FilterMetricsByLocation(metrics, tt.dc, tt.rack, tt.node) if len(filtered) != len(tt.wantIDs) { t.Fatalf("got %d metrics, want %d", len(filtered), len(tt.wantIDs)) }