/
publisher.go
161 lines (143 loc) · 4.94 KB
/
publisher.go
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
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
package publisher
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/Azure/go-shuttle/common/errorhandling"
"github.com/Azure/go-shuttle/marshal"
amqp "github.com/Azure/azure-amqp-common-go/v3"
servicebus "github.com/Azure/azure-service-bus-go"
"github.com/devigned/tab"
"github.com/Azure/go-shuttle/common"
"github.com/Azure/go-shuttle/internal/reflection"
"github.com/Azure/go-shuttle/prometheus/publisher"
)
type TopicPublisher interface {
common.Publisher
AppendTopicManagementOption(option servicebus.TopicManagementOption)
}
// Publisher is a struct to contain service bus entities relevant to publishing to a topic
type Publisher struct {
common.PublisherSettings
topic *servicebus.Topic
topicManagementOptions []servicebus.TopicManagementOption
}
func (p *Publisher) AppendTopicManagementOption(option servicebus.TopicManagementOption) {
p.topicManagementOptions = append(p.topicManagementOptions, option)
}
// New creates a new service bus publisher
func New(ctx context.Context, topicName string, opts ...ManagementOption) (*Publisher, error) {
ctx, s := tab.StartSpan(ctx, "go-shuttle.publisher.New")
defer s.End()
ns, err := servicebus.NewNamespace()
if err != nil {
return nil, err
}
publisher := &Publisher{PublisherSettings: common.PublisherSettings{}}
publisher.SetNamespace(ns)
publisher.SetMarshaller(marshal.JSONMarshaller)
for _, opt := range opts {
err := opt(publisher)
if err != nil {
return nil, err
}
}
topicEntity, err := ensureTopic(ctx, topicName, publisher.Namespace(), publisher.topicManagementOptions...)
if err != nil {
return nil, fmt.Errorf("failed to get topic: %w", err)
}
if err = publisher.initTopic(topicEntity.Name); err != nil {
return nil, fmt.Errorf("failed to create new topic %s: %w", topicEntity.Name, err)
}
return publisher, nil
}
// Publish publishes to the pre-configured Service Bus topic
func (p *Publisher) Publish(ctx context.Context, msg interface{}, opts ...Option) error {
ctx, s := tab.StartSpan(ctx, "go-shuttle.publisher.Publish")
defer s.End()
msgJSON, err := json.Marshal(msg)
// adding in user properties to enable filtering on listener side
sbMsg := servicebus.NewMessageFromString(string(msgJSON))
sbMsg.UserProperties = make(map[string]interface{})
sbMsg.UserProperties["type"] = reflection.GetType(msg)
// add in custom headers setup at initialization time
for headerName, headerKey := range p.Headers() {
val := reflection.GetReflectionValue(msg, headerKey)
if val != nil {
sbMsg.UserProperties[headerName] = val
}
}
for _, opt := range opts {
err := opt(sbMsg)
if err != nil {
return err
}
}
err = p.topic.Send(ctx, sbMsg)
if err == nil {
return nil
}
// recover + retry
if recErr := p.tryRecoverTopic(ctx, err); recErr != nil {
publisher.Metrics.IncConnectionRecoveryFailure(err)
return fmt.Errorf("failed to recover topic on send failure %s. recoveryError : %w, sendError: %s", p.topic.Name, recErr, err)
}
publisher.Metrics.IncConnectionRecoverySuccess(err)
if err = p.topic.Send(ctx, sbMsg); err != nil {
return fmt.Errorf("failed to send message to topic %s after recovery: %w", p.topic.Name, err)
}
return nil
}
func (p *Publisher) Close(ctx context.Context) error {
ctx, s := tab.StartSpan(ctx, "go-shuttle.publisher.Close")
defer s.End()
return p.topic.Close(ctx)
}
func (p *Publisher) tryRecoverTopic(ctx context.Context, sendError error) error {
ctx, s := tab.StartSpan(ctx, "go-shuttle.publisher.tryRecoverTopic", tab.StringAttribute("error", sendError.Error()))
defer s.End()
if errorhandling.IsConnectionDead(sendError) {
// re-create topic/sender
if err := p.initTopic(p.topic.Name); err != nil {
return fmt.Errorf("failed to init topic on recovery: %w", err)
}
publisher.Metrics.IncConnectionRecoverySuccess(sendError)
return nil
}
return fmt.Errorf("error is not identified as recoverable: %w", sendError)
}
func ensureTopic(ctx context.Context, name string, namespace *servicebus.Namespace, opts ...servicebus.TopicManagementOption) (*servicebus.TopicEntity, error) {
attempt := 1
tm := namespace.NewTopicManager()
ensure := func() (interface{}, error) {
ctx, span := tab.StartSpan(ctx, "go-shuttle.publisher.ensureTopic")
span.AddAttributes(tab.Int64Attribute("retry.attempt", int64(attempt)))
te, err := tm.Get(ctx, name)
if err == nil {
return te, nil
}
te, err = tm.Put(ctx, name, opts...)
if err != nil {
attempt++
tab.For(ctx).Error(err)
// let all errors be retryable for now. application only hit this once on topic creation.
return nil, amqp.Retryable(err.Error())
}
return te, nil
}
entity, err := amqp.Retry(5, 1*time.Second, ensure)
if err != nil {
tab.For(ctx).Error(err)
return nil, err
}
return entity.(*servicebus.TopicEntity), nil
}
func (p *Publisher) initTopic(name string) error {
topic, err := p.Namespace().NewTopic(name)
if err != nil {
return fmt.Errorf("failed to create new topic %s: %w", name, err)
}
p.topic = topic
return nil
}