mirror of
https://github.com/openziti/zrok.git
synced 2024-11-07 08:44:14 +01:00
simpler amqp sink approach (#351)
This commit is contained in:
parent
aabf695bec
commit
5a2f6a1f72
@ -35,34 +35,35 @@ type amqpSink struct {
|
||||
conn *amqp.Connection
|
||||
ch *amqp.Channel
|
||||
queue amqp.Queue
|
||||
join chan struct{}
|
||||
connected bool
|
||||
}
|
||||
|
||||
func newAmqpSink(cfg *AmqpSinkConfig) (*amqpSink, error) {
|
||||
return &amqpSink{
|
||||
as := &amqpSink{
|
||||
cfg: cfg,
|
||||
join: make(chan struct{}),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *amqpSink) Start() (join chan struct{}, err error) {
|
||||
logrus.Info("started")
|
||||
return s.join, nil
|
||||
}
|
||||
return as, nil
|
||||
}
|
||||
|
||||
func (s *amqpSink) Handle(event ZitiEventJson) error {
|
||||
if !s.connected {
|
||||
if err := s.connect(); err != nil {
|
||||
return err
|
||||
}
|
||||
logrus.Infof("connected to '%v'", s.cfg.Url)
|
||||
s.connected = true
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
logrus.Infof("pushing '%v'", event)
|
||||
return s.ch.PublishWithContext(ctx, "", s.queue.Name, false, false, amqp.Publishing{
|
||||
err := s.ch.PublishWithContext(ctx, "", s.queue.Name, false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
Body: []byte(event),
|
||||
})
|
||||
}
|
||||
|
||||
func (s *amqpSink) Stop() {
|
||||
close(s.join)
|
||||
logrus.Info("stopped")
|
||||
if err != nil {
|
||||
s.connected = false
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *amqpSink) connect() (err error) {
|
||||
@ -70,18 +71,13 @@ func (s *amqpSink) connect() (err error) {
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "error dialing amqp broker")
|
||||
}
|
||||
|
||||
s.ch, err = s.conn.Channel()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "error getting amqp channel")
|
||||
}
|
||||
|
||||
s.queue, err = s.ch.QueueDeclare(s.cfg.QueueName, true, false, false, false, nil)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "error declaring queue")
|
||||
}
|
||||
|
||||
logrus.Infof("connected to amqp broker at '%v'", s.cfg.Url)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
@ -14,7 +14,6 @@ type Bridge struct {
|
||||
src ZitiEventJsonSource
|
||||
srcJoin chan struct{}
|
||||
snk ZitiEventJsonSink
|
||||
snkJoin chan struct{}
|
||||
events chan ZitiEventMsg
|
||||
close chan struct{}
|
||||
join chan struct{}
|
||||
@ -40,9 +39,6 @@ func NewBridge(cfg *BridgeConfig) (*Bridge, error) {
|
||||
}
|
||||
|
||||
func (b *Bridge) Start() (join chan struct{}, err error) {
|
||||
if b.snkJoin, err = b.snk.Start(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if b.srcJoin, err = b.src.Start(b.events); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@ -76,9 +72,7 @@ func (b *Bridge) Start() (join chan struct{}, err error) {
|
||||
|
||||
func (b *Bridge) Stop() {
|
||||
b.src.Stop()
|
||||
b.snk.Stop()
|
||||
close(b.close)
|
||||
<-b.srcJoin
|
||||
<-b.snkJoin
|
||||
<-b.join
|
||||
}
|
||||
|
@ -79,7 +79,5 @@ type ZitiEventJsonSource interface {
|
||||
}
|
||||
|
||||
type ZitiEventJsonSink interface {
|
||||
Start() (join chan struct{}, err error)
|
||||
Handle(event ZitiEventJson) error
|
||||
Stop()
|
||||
}
|
||||
|
Loading…
Reference in New Issue
Block a user