zrok/controller/metrics/amqpSink.go

82 lines
2.0 KiB
Go
Raw Normal View History

2023-03-15 21:14:06 +01:00
package metrics
2023-03-15 17:47:26 +01:00
import (
"context"
"github.com/michaelquigley/cf"
"github.com/openziti/zrok/controller/env"
"github.com/pkg/errors"
amqp "github.com/rabbitmq/amqp091-go"
2023-04-05 19:57:22 +02:00
"github.com/sirupsen/logrus"
2023-03-15 17:47:26 +01:00
"time"
)
func init() {
env.GetCfOptions().AddFlexibleSetter("amqpSink", loadAmqpSinkConfig)
}
type AmqpSinkConfig struct {
Url string `cf:"+secret"`
QueueName string
}
func loadAmqpSinkConfig(v interface{}, _ *cf.Options) (interface{}, error) {
if submap, ok := v.(map[string]interface{}); ok {
cfg := &AmqpSinkConfig{}
if err := cf.Bind(cfg, submap, cf.DefaultOptions()); err != nil {
return nil, err
}
return newAmqpSink(cfg)
}
return nil, errors.New("invalid config structure for 'amqpSink'")
}
type amqpSink struct {
2023-06-21 17:33:43 +02:00
cfg *AmqpSinkConfig
conn *amqp.Connection
ch *amqp.Channel
queue amqp.Queue
connected bool
2023-03-15 17:47:26 +01:00
}
func newAmqpSink(cfg *AmqpSinkConfig) (*amqpSink, error) {
2023-06-21 17:41:13 +02:00
as := &amqpSink{cfg: cfg}
2023-06-21 17:33:43 +02:00
return as, nil
2023-03-15 17:47:26 +01:00
}
func (s *amqpSink) Handle(event ZitiEventJson) error {
2023-06-21 17:33:43 +02:00
if !s.connected {
if err := s.connect(); err != nil {
return err
}
logrus.Infof("connected to '%v'", s.cfg.Url)
s.connected = true
}
2023-03-15 17:47:26 +01:00
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
2023-04-05 19:57:22 +02:00
logrus.Infof("pushing '%v'", event)
2023-06-21 17:33:43 +02:00
err := s.ch.PublishWithContext(ctx, "", s.queue.Name, false, false, amqp.Publishing{
2023-03-15 17:47:26 +01:00
ContentType: "application/json",
Body: []byte(event),
})
2023-06-21 17:33:43 +02:00
if err != nil {
s.connected = false
}
return err
}
func (s *amqpSink) connect() (err error) {
s.conn, err = amqp.Dial(s.cfg.Url)
if err != nil {
return errors.Wrapf(err, "error dialing '%v'", s.cfg.Url)
}
s.ch, err = s.conn.Channel()
if err != nil {
return errors.Wrapf(err, "error getting amqp channel from '%v'", s.cfg.Url)
}
s.queue, err = s.ch.QueueDeclare(s.cfg.QueueName, true, false, false, false, nil)
if err != nil {
return errors.Wrapf(err, "error declaring queue '%v' with '%v'", s.cfg.QueueName, s.cfg.Url)
}
return nil
}