131
hostInfo membership.HostInfo,
132
serializer serialization.Serializer,
134
>
return &Scanner{
135
>
context: scannerContext{
136
>
cfg: cfg,
137
>
sdkClientFactory: sdkClientFactory,
138
>
logger: logger,
139
>
metricsHandler: metricsHandler,
140
>
executionManager: executionManager,
141
>
taskManager: taskManager,
142
>
visibilityManager: visibilityManager,
143
>
metadataManager: metadataManager,
144
>
historyClient: historyClient,
145
>
matchingClient: matchingClient,
146
>
adminClient: adminClient,
147
>
namespaceRegistry: registry,
148
>
currentClusterName: currentClusterName,
149
>
hostInfo: hostInfo,
150
>
serializer: serializer,
151
>
},
152
>
}
153
>
}
154
155
// Start starts the scanner
156
>
func (s *Scanner) Start() error {
scanner.go
157
>
ctx := context.WithValue(context.Background(), scannerContextKey, s.context)
158
>
ctx = headers.SetCallerInfo(ctx, headers.SystemBackgroundHighCallerInfo)
159
>
ctx, s.lifecycleCancel = context.WithCancel(ctx)
160
>
161
>
workerOpts := worker.Options{
162
>
Identity: "temporal-system@" + s.context.hostInfo.Identity(),
163
>
MaxConcurrentActivityExecutionSize: s.context.cfg.MaxConcurrentActivityExecutionSize(),
164
>
MaxConcurrentWorkflowTaskExecutionSize: s.context.cfg.MaxConcurrentWorkflowTaskExecutionSize(),
165
>
MaxConcurrentActivityTaskPollers: s.context.cfg.MaxConcurrentActivityTaskPollers(),
166
>
MaxConcurrentWorkflowTaskPollers: s.context.cfg.MaxConcurrentWorkflowTaskPollers(),
167
>
168
>
BackgroundActivityContext: ctx,
169
>
}
170
>
171
>
var workerTaskQueueNames []string
172
>
if s.context.cfg.Persistence.DefaultStoreType() != config.StoreTypeSQL && s.context.cfg.ExecutionsScannerEnabled() {
173
s.wg.Add(1)
174
go s.startWorkflowWithRetry(ctx, executionsScannerWFStartOptions, executionsScannerWFTypeName)
175
workerTaskQueueNames = append(workerTaskQueueNames, executionsScannerTaskQueueName)
176
>
} else if s.context.cfg.ExecutionsScannerEnabled() {
scanner.go
177
s.context.logger.Info("ExecutionsScanner is not supported for SQL store")
178
}
179
180
>
if s.context.cfg.Persistence.DefaultStoreType() == config.StoreTypeSQL && s.context.cfg.TaskQueueScannerEnabled() {
scanner.go
181
s.wg.Add(1)
182
go s.startWorkflowWithRetry(ctx, tlScannerWFStartOptions, tqScannerWFTypeName)