diff --git a/workers/go/harness/worker.go b/workers/go/harness/worker.go index 98423583..2963920f 100644 --- a/workers/go/harness/worker.go +++ b/workers/go/harness/worker.go @@ -168,7 +168,6 @@ func buildWorkerOptions(flags *pflag.FlagSet, args clioptions.WorkerOptions, pro } if flags.Changed("activity-poller-autoscale-max") { options.ActivityTaskPollerBehavior = sdkworker.NewPollerBehaviorAutoscaling(sdkworker.PollerBehaviorAutoscalingOptions{ - InitialNumberOfPollers: args.ActivityPollerAutoscaleMax, MaximumNumberOfPollers: args.ActivityPollerAutoscaleMax, }) } else if flags.Changed("max-concurrent-activity-pollers") { @@ -178,7 +177,6 @@ func buildWorkerOptions(flags *pflag.FlagSet, args clioptions.WorkerOptions, pro } if flags.Changed("workflow-poller-autoscale-max") { options.WorkflowTaskPollerBehavior = sdkworker.NewPollerBehaviorAutoscaling(sdkworker.PollerBehaviorAutoscalingOptions{ - InitialNumberOfPollers: args.WorkflowPollerAutoscaleMax, MaximumNumberOfPollers: args.WorkflowPollerAutoscaleMax, }) } else if flags.Changed("max-concurrent-workflow-pollers") { diff --git a/workers/go/harness/worker_test.go b/workers/go/harness/worker_test.go index ed5756e2..5f355cc1 100644 --- a/workers/go/harness/worker_test.go +++ b/workers/go/harness/worker_test.go @@ -3,9 +3,11 @@ package harness import ( "context" "errors" + "reflect" "testing" "time" + "github.com/temporalio/omes/clioptions" sdkclient "go.temporal.io/sdk/client" sdkworker "go.temporal.io/sdk/worker" ) @@ -94,3 +96,60 @@ func TestRunWorkersStopsRemainingWorkersOnFailure(t *testing.T) { t.Fatal("expected remaining worker to observe shutdown") } } + +func TestBuildWorkerOptionsPollerBehavior(t *testing.T) { + for _, tc := range []struct { + name string + argv []string + activityBehavior sdkworker.PollerBehavior + workflowBehavior sdkworker.PollerBehavior + }{ + { + name: "no poller flags leaves both unset", + argv: nil, + }, + { + // The autoscaling max must not also pin the initial poller count; + // unset fields fall back to the SDK's own defaults. + name: "activity autoscale max sets only the maximum", + argv: []string{"--activity-poller-autoscale-max=20"}, + activityBehavior: sdkworker.NewPollerBehaviorAutoscaling(sdkworker.PollerBehaviorAutoscalingOptions{ + MaximumNumberOfPollers: 20, + }), + }, + { + name: "workflow autoscale max sets only the maximum", + argv: []string{"--workflow-poller-autoscale-max=10"}, + workflowBehavior: sdkworker.NewPollerBehaviorAutoscaling(sdkworker.PollerBehaviorAutoscalingOptions{ + MaximumNumberOfPollers: 10, + }), + }, + { + name: "max concurrent pollers still pins a simple maximum", + argv: []string{"--max-concurrent-activity-pollers=3"}, + activityBehavior: sdkworker.NewPollerBehaviorSimpleMaximum(sdkworker.PollerBehaviorSimpleMaximumOptions{ + MaximumNumberOfPollers: 3, + }), + }, + } { + t.Run(tc.name, func(t *testing.T) { + var args clioptions.WorkerOptions + flags := args.FlagSetWithPrefix("") + if err := flags.Parse(tc.argv); err != nil { + t.Fatalf("failed to parse %v: %v", tc.argv, err) + } + options, err := buildWorkerOptions(flags, args, "") + if err != nil { + t.Fatalf("buildWorkerOptions failed: %v", err) + } + if !reflect.DeepEqual(options.ActivityTaskPollerBehavior, tc.activityBehavior) { + t.Errorf("activity poller behavior = %#v, want %#v", + options.ActivityTaskPollerBehavior, tc.activityBehavior) + } + if !reflect.DeepEqual(options.WorkflowTaskPollerBehavior, tc.workflowBehavior) { + t.Errorf("workflow poller behavior = %#v, want %#v", + options.WorkflowTaskPollerBehavior, tc.workflowBehavior) + } + }) + } +}