471
components []workercommon.PerNSWorkerComponent,
472
allocation workerAllocation,
474
>
nsName := w.ns.Name().String()
475
>
// this should not block because it uses an existing grpc connection
476
>
client := w.wm.sdkClientFactory.NewClient(sdkclient.Options{
477
>
Namespace: nsName,
478
>
DataConverter: sdk.PreferProtoDataConverter,
479
>
})
480
>
481
>
var sdkoptions sdkworker.Options
482
>
483
>
// copy from dynamic config. apply explicit defaults for some instead of using the sdk
484
>
// defaults so that we can multiply below.
485
>
sdkoptions.MaxConcurrentActivityExecutionSize = cmp.Or(w.opts.MaxConcurrentActivityExecutionSize, 1000)
486
>
sdkoptions.WorkerActivitiesPerSecond = w.opts.WorkerActivitiesPerSecond
487
>
sdkoptions.MaxConcurrentLocalActivityExecutionSize = cmp.Or(w.opts.MaxConcurrentLocalActivityExecutionSize, 1000)
488
>
sdkoptions.WorkerLocalActivitiesPerSecond = w.opts.WorkerLocalActivitiesPerSecond
489
>
sdkoptions.MaxConcurrentActivityTaskPollers = max(cmp.Or(w.opts.MaxConcurrentActivityTaskPollers, 2), 2)
490
>
sdkoptions.MaxConcurrentWorkflowTaskExecutionSize = cmp.Or(w.opts.MaxConcurrentWorkflowTaskExecutionSize, 1000)
491
>
sdkoptions.MaxConcurrentWorkflowTaskPollers = max(cmp.Or(w.opts.MaxConcurrentWorkflowTaskPollers, 2), 2)
492
>
sdkoptions.StickyScheduleToStartTimeout = w.opts.StickyScheduleToStartTimeout
493
>
494
>
sdkoptions.BackgroundActivityContext = headers.SetCallerInfo(context.Background(), headers.NewBackgroundHighCallerInfo(nsName))
495
>
sdkoptions.Identity = fmt.Sprintf("temporal-system@%s@%s", w.wm.hostName, nsName)
496
>
// increase these if we're supposed to run with more allocation
497
>
sdkoptions.MaxConcurrentWorkflowTaskPollers *= allocation.local
498
>
sdkoptions.MaxConcurrentActivityTaskPollers *= allocation.local
499
>
sdkoptions.MaxConcurrentLocalActivityExecutionSize *= allocation.local
500
>
sdkoptions.MaxConcurrentWorkflowTaskExecutionSize *= allocation.local
501
>
sdkoptions.MaxConcurrentActivityExecutionSize *= allocation.local
502
>
sdkoptions.OnFatalError = w.onFatalError
503
>
504
>
// this should not block because the client already has server capabilities
505
>
worker := w.wm.sdkClientFactory.NewWorker(client, w.wm.taskQueueName, sdkoptions)
506
>
details := workercommon.RegistrationDetails{
507
>
TotalWorkers: allocation.total,
508
>
Multiplicity: allocation.local,
509
>
}
510
>
for _, cmp := range components {
511
>
cleanup := cmp.Register(worker, w.ns, details)
512
>
if cleanup != nil {
513
w.cleanup = append(w.cleanup, cleanup)
514
}