105
testHooks testhooks.TestHooks,
106
dlqWriter DLQWriter,
108
>
shardID := shardContext.GetShardID()
109
>
taskRetryPolicy := backoff.NewExponentialRetryPolicy(config.ReplicationTaskProcessorErrorRetryWait(shardID)).
110
>
WithBackoffCoefficient(config.ReplicationTaskProcessorErrorRetryBackoffCoefficient(shardID)).
111
>
WithMaximumInterval(config.ReplicationTaskProcessorErrorRetryMaxInterval(shardID)).
112
>
WithMaximumAttempts(config.ReplicationTaskProcessorErrorRetryMaxAttempts(shardID)).
113
>
WithExpirationInterval(config.ReplicationTaskProcessorErrorRetryExpiration(shardID))
114
>
115
>
// TODO: define separate set of configs for dlq retry
116
>
dlqRetryPolicy := backoff.NewExponentialRetryPolicy(config.ReplicationTaskProcessorErrorRetryWait(shardID)).
117
>
WithBackoffCoefficient(config.ReplicationTaskProcessorErrorRetryBackoffCoefficient(shardID)).
118
>
WithMaximumInterval(config.ReplicationTaskProcessorErrorRetryMaxInterval(shardID)).
119
>
WithMaximumAttempts(config.ReplicationTaskProcessorErrorRetryMaxAttempts(shardID)).
120
>
WithExpirationInterval(config.ReplicationTaskProcessorErrorRetryExpiration(shardID))
121
>
122
>
return &taskProcessorImpl{
123
>
status: common.DaemonStatusInitialized,
124
>
sourceShardID: sourceShardID,
125
>
sourceCluster: replicationTaskFetcher.getSourceCluster(),
126
>
shard: shardContext,
127
>
historyEngine: historyEngine,
128
>
historySerializer: eventSerializer,
129
>
config: config,
130
>
metricsHandler: metricsHandler,
131
>
logger: shardContext.GetLogger(),
132
>
testHooks: testHooks,
133
>
replicationTaskExecutor: replicationTaskExecutor,
134
>
dlqWriter: dlqWriter,
135
>
rateLimiter: quotas.NewMultiRateLimiter([]quotas.RateLimiter{
136
>
quotas.NewDefaultOutgoingRateLimiter(
137
>
func() float64 { return config.ReplicationTaskProcessorShardQPS() },
138
),
139
replicationTaskFetcher.getRateLimiter(),