-
Notifications
You must be signed in to change notification settings - Fork 0
/
publish.go
executable file
·154 lines (123 loc) · 3.26 KB
/
publish.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
package rabbitmqv1
import (
"encoding/json"
"github.com/streadway/amqp"
"golang.org/x/net/context"
"reflect"
"strconv"
"sync"
"time"
)
var (
deliveryMode uint8 = 2
headerError = "Error"
headerRetryCount = "RetryCount"
headerStackTrace = "StackTrace"
headerTime = "Time"
maxChannelMutex = &sync.Mutex{}
)
const resendDelay = 1 * time.Millisecond
type (
PublisherBuilder struct {
openChannelCount int
messageBroker MessageBroker
brokerChannel *Channel
publishers []Publisher
startPublisherCh chan bool
isAlreadyStartConnection bool
isChannelActive bool
}
Publisher struct {
isAlreadyCreated bool
exchangeName string
exchangeType ExchangeType
payloadTypes []reflect.Type
}
)
func errorPublishMessage(correlationId string, payload []byte, retryCount int, err error, stackTracing string) amqp.Publishing {
headers := make(map[string]interface{})
headers[headerRetryCount] = strconv.Itoa(retryCount)
headers[headerError] = err.Error()
headers[headerStackTrace] = stackTracing
headers[headerTime] = time.Now().String()
return amqp.Publishing{
MessageId: getGuid(),
Body: payload,
Headers: headers,
CorrelationId: correlationId,
Timestamp: time.Now(),
DeliveryMode: deliveryMode,
ContentEncoding: "UTF-8",
ContentType: "use-cases/json",
}
}
func getBytes(key interface{}) ([]byte, error) {
return json.Marshal(key)
}
func (p *PublisherBuilder) SubscriberExchange() *PublisherBuilder {
for index, item := range p.publishers {
if !item.isAlreadyCreated {
p.brokerChannel.createExchange(item.exchangeName, item.exchangeType, nil)
p.publishers[index].isAlreadyCreated = true
}
}
return p
}
func (p *PublisherBuilder) Publish(ctx context.Context, routingKey string, exchangeName string, payload interface{}) error {
var (
err error
message amqp.Publishing
)
for {
if p.isChannelActive {
break
}
}
p.SubscriberExchange()
if message, err = publishMessage("", payload); err != nil {
return err
}
err = p.brokerChannel.rabbitChannel.Publish(exchangeName, routingKey, false, false, message)
select {
case confirm := <-p.brokerChannel.notifyConfirm:
if confirm.Ack {
break
}
case <-time.After(resendDelay):
}
return err
}
func publishMessage(correlationId string, payload interface{}) (amqp.Publishing, error) {
headers := make(map[string]interface{})
headers[headerTime] = time.Now().String()
bodyJson, err := json.Marshal(payload)
return amqp.Publishing{
MessageId: getGuid(),
Body: bodyJson,
Headers: headers,
CorrelationId: correlationId,
Timestamp: time.Now(),
DeliveryMode: deliveryMode,
ContentEncoding: "UTF-8",
ContentType: "use-cases/json",
}, err
}
func (p *PublisherBuilder) CreateChannel() {
var err error
go func() {
for {
select {
case isConnected := <-p.startPublisherCh:
if isConnected {
if p.brokerChannel, err = p.messageBroker.CreateChannel(); err != nil {
panic(err)
}
p.isChannelActive = true
p.brokerChannel.rabbitChannel.NotifyPublish(p.brokerChannel.notifyConfirm)
} else {
p.isChannelActive = false
}
}
}
}()
}