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