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")
scanner.go
178
>
}
179
180
>
if s.context.cfg.Persistence.DefaultStoreType() == config.StoreTypeSQL && s.context.cfg.TaskQueueScannerEnabled() {
scanner.go
182
>
go s.startWorkflowWithRetry(ctx, tlScannerWFStartOptions, tqScannerWFTypeName)
183
>
workerTaskQueueNames = append(workerTaskQueueNames, tqScannerTaskQueueName)
184
>
}
185
186
>
if s.context.cfg.HistoryScannerEnabled() {
scanner.go
188
>
go s.startWorkflowWithRetry(ctx, historyScannerWFStartOptions, historyScannerWFTypeName)
189
>
workerTaskQueueNames = append(workerTaskQueueNames, historyScannerTaskQueueName)
190
>
}
191
192
>
if s.context.cfg.BuildIdScavengerEnabled() {
scanner.go
194
>
go s.startWorkflowWithRetry(ctx, build_ids.BuildIdScavengerWFStartOptions, build_ids.BuildIdScavangerWorkflowName)
195
>
196
>
buildIdsActivities := build_ids.NewActivities(
197
>
s.context.logger,
198
>
s.context.taskManager,
199
>
s.context.metadataManager,
200
>
s.context.visibilityManager,
201
>
s.context.namespaceRegistry,
202
>
s.context.matchingClient,
203
>
s.context.currentClusterName,
204
>
s.context.cfg.RemovableBuildIdDurationSinceDefault,
205
>
s.context.cfg.BuildIdScavengerVisibilityRPS,
206
>
)
207
>
208
>
work := s.context.sdkClientFactory.NewWorker(s.context.sdkClientFactory.GetSystemClient(), build_ids.BuildIdScavengerTaskQueueName, workerOpts)
209
>
work.RegisterWorkflowWithOptions(build_ids.BuildIdScavangerWorkflow, workflow.RegisterOptions{Name: build_ids.BuildIdScavangerWorkflowName})
210
>
work.RegisterActivityWithOptions(buildIdsActivities.ScavengeBuildIds, activity.RegisterOptions{Name: build_ids.BuildIdScavangerActivityName})
211
>
212
>
// TODO: Nothing is gracefully stopping these workers or listening for fatal errors.
213
>
if err := work.Start(); err != nil {
214
return err
215
}
216
}
217
218
>
siOpts := s.context.cfg.ScheduleInvariantsScannerOptions()
scanner.go
219
>
if siOpts.OverdueNextActionTimeEnabled || siOpts.StuckOpenEnabled || siOpts.UnknownStateEnabled {
220
scheduleActivities := scheduleinvariants.NewActivities(
221
s.context.logger,