Skip to content
File

Blob: spec/chord/retry.go

go108 lines
1package chord
2 
3import (
4 "context"
5 "expvar"
6 "time"
7 
8 "go.miragespace.co/specter/spec/protocol"
9 
10 "github.com/avast/retry-go/v5"
11)
12 
13var kvRetries = expvar.NewInt("chord.kvRetries")
14 
15type retryableWrapper struct {
16 VNode
17 retryInterval time.Duration
18 retryAttempts uint
19}
20 
21// WrapRetryKV wraps a given VNode to provide automatic retry on retryable KV errors
22func WrapRetryKV(vnode VNode, interval time.Duration, maxAttempts uint) VNode {
23 return &retryableWrapper{
24 VNode: vnode,
25 retryInterval: interval,
26 retryAttempts: maxAttempts,
27 }
28}
29 
30func (n *retryableWrapper) retryOptions(ctx context.Context) []retry.Option {
31 return []retry.Option{
32 retry.Context(ctx),
33 retry.Attempts(n.retryAttempts),
34 retry.Delay(n.retryInterval),
35 retry.OnRetry(func(n uint, err error) {
36 kvRetries.Add(1)
37 }),
38 retry.RetryIf(ErrorIsRetryable),
39 retry.LastErrorOnly(true),
40 }
41}
42 
43func (n *retryableWrapper) Put(ctx context.Context, key []byte, value []byte) error {
44 return retry.New(n.retryOptions(ctx)...).Do(func() error {
45 return n.VNode.Put(ctx, key, value)
46 })
47}
48 
49func (n *retryableWrapper) Get(ctx context.Context, key []byte) (value []byte, err error) {
50 return retry.NewWithData[[]byte](n.retryOptions(ctx)...).Do(func() ([]byte, error) {
51 return n.VNode.Get(ctx, key)
52 })
53}
54 
55func (n *retryableWrapper) Delete(ctx context.Context, key []byte) error {
56 return retry.New(n.retryOptions(ctx)...).Do(func() error {
57 return n.VNode.Delete(ctx, key)
58 })
59}
60 
61func (n *retryableWrapper) PrefixAppend(ctx context.Context, prefix []byte, child []byte) error {
62 return retry.New(n.retryOptions(ctx)...).Do(func() error {
63 return n.VNode.PrefixAppend(ctx, prefix, child)
64 })
65}
66 
67func (n *retryableWrapper) PrefixList(ctx context.Context, prefix []byte) (children [][]byte, err error) {
68 return retry.NewWithData[[][]byte](n.retryOptions(ctx)...).Do(func() ([][]byte, error) {
69 return n.VNode.PrefixList(ctx, prefix)
70 })
71}
72 
73func (n *retryableWrapper) PrefixContains(ctx context.Context, prefix []byte, child []byte) (bool, error) {
74 return retry.NewWithData[bool](n.retryOptions(ctx)...).Do(func() (bool, error) {
75 return n.VNode.PrefixContains(ctx, prefix, child)
76 })
77}
78 
79func (n *retryableWrapper) PrefixRemove(ctx context.Context, prefix []byte, child []byte) error {
80 return retry.New(n.retryOptions(ctx)...).Do(func() error {
81 return n.VNode.PrefixRemove(ctx, prefix, child)
82 })
83}
84 
85func (n *retryableWrapper) Acquire(ctx context.Context, lease []byte, ttl time.Duration) (token uint64, err error) {
86 return retry.NewWithData[uint64](n.retryOptions(ctx)...).Do(func() (uint64, error) {
87 return n.VNode.Acquire(ctx, lease, ttl)
88 })
89}
90 
91func (n *retryableWrapper) Renew(ctx context.Context, lease []byte, ttl time.Duration, prevToken uint64) (newToken uint64, err error) {
92 return retry.NewWithData[uint64](n.retryOptions(ctx)...).Do(func() (uint64, error) {
93 return n.VNode.Renew(ctx, lease, ttl, prevToken)
94 })
95}
96 
97func (n *retryableWrapper) Release(ctx context.Context, lease []byte, token uint64) error {
98 return retry.New(n.retryOptions(ctx)...).Do(func() error {
99 return n.VNode.Release(ctx, lease, token)
100 })
101}
102 
103func (n *retryableWrapper) ListKeys(ctx context.Context, prefix []byte) ([]*protocol.KeyComposite, error) {
104 return retry.NewWithData[[]*protocol.KeyComposite](n.retryOptions(ctx)...).Do(func() ([]*protocol.KeyComposite, error) {
105 return n.VNode.ListKeys(ctx, prefix)
106 })
107}