-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathrequeue_retry.go
More file actions
79 lines (67 loc) · 2.03 KB
/
Copy pathrequeue_retry.go
File metadata and controls
79 lines (67 loc) · 2.03 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
package httpsqs
import (
"context"
"time"
"github.com/gomooth/httpsqs"
"github.com/gomooth/pkg/framework/retry"
"github.com/gomooth/pkg/mq/internal/attempt_tracker"
mqretry "github.com/gomooth/pkg/mq/internal/retry"
"github.com/gomooth/pkg/mq/internal/logutil"
"github.com/gomooth/pkg/mq/internal/metrics"
"github.com/gomooth/pkg/mq/internal/types"
)
// requeueRetryStrategy 再入队重试策略,内部委托给 mqretry.RequeueStrategy
type requeueRetryStrategy struct {
inner *mqretry.RequeueStrategy
handler types.IHandler
}
func newRequeueRetryStrategy(
handler types.IHandler,
maxRetry int,
backoff retry.BackoffStrategy,
client httpsqs.IClient,
queueName string,
_ logutil.Logger,
m *metrics.ConsumerMetrics,
) *requeueRetryStrategy {
backoffFn := mqretry.BackoffDelayFunc(func(attempt uint) time.Duration {
return backoff.Delay(attempt)
})
tracker := attempt_tracker.NewAttemptTracker()
requeueFn := func(ctx context.Context, msg types.Message) error {
_, pushErr := client.Put(ctx, queueName, string(msg.Data))
return pushErr
}
return &requeueRetryStrategy{
handler: handler,
inner: mqretry.NewRequeueStrategy(mqretry.RequeueConfig{
MaxRetry: maxRetry,
Backoff: backoffFn,
Tracker: tracker,
Metrics: m,
Requeue: requeueFn,
}),
}
}
func (s *requeueRetryStrategy) SetFailedHandler(fn types.FailedHandlerFunc) {
s.inner.SetFailedHandler(fn)
}
func (s *requeueRetryStrategy) SetDeadLetterHandler(h types.DeadLetterHandler) {
s.inner.SetDeadLetterHandler(h)
}
func (s *requeueRetryStrategy) SetTimeout(d time.Duration) {
s.inner.SetTimeout(d)
}
func (s *requeueRetryStrategy) OnMessage(ctx context.Context, queue string, data []byte) error {
msg := types.NewHttpsqSMessage(queue, data, 0)
err := s.inner.OnMessage(ctx, msg, s.handler.Handle)
// 兼容旧行为:上下文取消时返回错误,其他情况返回 nil
if err != nil && ctx.Err() != nil {
return err
}
return nil
}
// Close 停止 AttemptTracker 的后台清理 goroutine
func (s *requeueRetryStrategy) Close() {
s.inner.Close()
}