86
enableDataLossMetrics EnableDataLossMetrics,
87
enableBestEffortDeleteTasksOnWorkflowUpdate EnableBestEffortDeleteTasksOnWorkflowUpdate,
89
>
factory := &factoryImpl{
90
>
dataStoreFactory: dataStoreFactory,
91
>
config: cfg,
92
>
serializer: serializer,
93
>
eventBlobCache: eventBlobCache,
94
>
metricsHandler: metricsHandler,
95
>
logger: logger,
96
>
clusterName: clusterName,
97
>
systemRateLimiter: systemRateLimiter,
98
>
namespaceRateLimiter: namespaceRateLimiter,
99
>
shardRateLimiter: shardRateLimiter,
100
>
healthSignals: healthSignals,
101
>
enableDataLossMetrics: dynamicconfig.BoolPropertyFn(enableDataLossMetrics),
102
>
enableBestEffortDeleteTasksOnWorkflowUpdate: dynamicconfig.BoolPropertyFn(enableBestEffortDeleteTasksOnWorkflowUpdate),
103
>
}
104
>
factory.initDependencies()
105
>
return factory
106
>
}
107
108
// NewTaskManager returns a new task manager
109
>
func (f *factoryImpl) NewTaskManager() (persistence.TaskManager, error) {
factory.go
110
>
taskStore, err := f.dataStoreFactory.NewTaskStore()
111
>
if err != nil {
112
return nil, err
113
}
114
>
result := persistence.NewTaskManager(taskStore, f.serializer)
factory.go
115
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
116
result = persistence.NewTaskPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
117
}
118
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
119
>
result = persistence.NewTaskPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
120
>
}
121
>
result = persistence.NewTaskPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
122
>
return result, nil
123
}
124
125
// NewFairTaskManager returns a new task fairness manager
126
>
func (f *factoryImpl) NewFairTaskManager() (persistence.FairTaskManager, error) {
factory.go
127
>
taskStore, err := f.dataStoreFactory.NewFairTaskStore()
128
>
if err != nil {
129
return nil, err
130
}
131
>
result := persistence.NewTaskManager(taskStore, f.serializer)
factory.go
132
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
133
result = persistence.NewTaskPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
134
}
135
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
136
>
result = persistence.NewTaskPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
137
>
}
138
>
result = persistence.NewTaskPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
139
>
return result, nil
140
}
141
142
// NewShardManager returns a new shard manager
143
>
func (f *factoryImpl) NewShardManager() (persistence.ShardManager, error) {
factory.go
144
>
shardStore, err := f.dataStoreFactory.NewShardStore()
145
>
if err != nil {
146
return nil, err
147
}
148
149
>
result := persistence.NewShardManager(shardStore, f.serializer)
factory.go
150
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
151
result = persistence.NewShardPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
152
}
153
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
154
>
result = persistence.NewShardPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
155
>
}
156
>
result = persistence.NewShardPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
157
>
return result, nil
158
}
159
160
// NewMetadataManager returns a new metadata manager
161
>
func (f *factoryImpl) NewMetadataManager() (persistence.MetadataManager, error) {
factory.go
162
>
store, err := f.dataStoreFactory.NewMetadataStore()
163
>
if err != nil {
164
return nil, err
165
}
166
167
>
result := persistence.NewMetadataManagerImpl(store, f.serializer, f.logger, f.clusterName)
factory.go
168
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
169
result = persistence.NewMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
170
}
171
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
172
>
result = persistence.NewMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
173
>
}
174
>
result = persistence.NewMetadataPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
175
>
return result, nil
176
}
177
178
// NewClusterMetadataManager returns a new cluster metadata manager
179
>
func (f *factoryImpl) NewClusterMetadataManager() (persistence.ClusterMetadataManager, error) {
factory.go
180
>
store, err := f.dataStoreFactory.NewClusterMetadataStore()
181
>
if err != nil {
182
return nil, err
183
}
184
185
>
result := persistence.NewClusterMetadataManagerImpl(store, f.serializer, f.clusterName, f.logger)
factory.go
186
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
187
result = persistence.NewClusterMetadataPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
188
}
189
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
190
>
result = persistence.NewClusterMetadataPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
191
>
}
192
>
result = persistence.NewClusterMetadataPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
193
>
return result, nil
194
}
195
196
// NewExecutionManager returns a new execution manager
197
>
func (f *factoryImpl) NewExecutionManager() (persistence.ExecutionManager, error) {
factory.go
198
>
store, err := f.dataStoreFactory.NewExecutionStore()
199
>
if err != nil {
200
return nil, err
201
}
202
203
>
result := persistence.NewExecutionManager(
factory.go
204
>
store,
205
>
f.serializer,
206
>
f.eventBlobCache,
207
>
f.logger,
208
>
f.config.TransactionSizeLimit,
209
>
f.enableBestEffortDeleteTasksOnWorkflowUpdate,
210
>
)
211
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
212
result = persistence.NewExecutionPersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
213
}
214
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
215
>
result = persistence.NewExecutionPersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
216
>
}
217
>
result = persistence.NewExecutionPersistenceRetryableClient(result, retryPolicy, IsPersistenceTransientError)
218
>
return result, nil
219
}
220
221
>
func (f *factoryImpl) NewNamespaceReplicationQueue() (persistence.NamespaceReplicationQueue, error) {
factory.go
222
>
result, err := f.dataStoreFactory.NewQueue(persistence.NamespaceReplicationQueueType)
223
>
if err != nil {
224
return nil, err
225
}
226
227
>
if f.systemRateLimiter != nil && f.namespaceRateLimiter != nil {
factory.go
228
result = persistence.NewQueuePersistenceRateLimitedClient(result, f.systemRateLimiter, f.namespaceRateLimiter, f.shardRateLimiter, f.logger)
229
}
230
>
if f.metricsHandler != nil && f.healthSignals != nil {
factory.go
231
>
result = persistence.NewQueuePersistenceMetricsClient(result, f.metricsHandler, f.healthSignals, f.logger, f.enableDataLossMetrics)
232
>
}
233
>
result = persistence.NewQueuePersistenceRetryableClient(result, namespaceQueueRetryPolicy, IsNamespaceQueueTransientError)
234
>
return persistence.NewNamespaceReplicationQueue(result, f.serializer, f.clusterName, f.metricsHandler, f.logger)
235
}
236