Add amqp support

This commit is contained in:
2023-10-14 19:37:39 +03:00
parent 7b7c419f03
commit 234a50adb2
8 changed files with 168 additions and 11 deletions
+79
View File
@@ -0,0 +1,79 @@
package ui
import (
"cli-mon/yabl"
"context"
json "encoding/json"
"fmt"
amqp "github.com/rabbitmq/amqp091-go"
"os"
"time"
)
const exchange = "yabl-newproto-in"
var hostname, _ = os.Hostname()
type broker struct {
uri string
conn *amqp.Connection
ch *amqp.Channel
}
func NewAmqp(uri string) UI {
return &broker{uri: uri}
}
func (b *broker) Consume(event *yabl.Event) {
var bytes []byte
var err error
if bytes, err = json.Marshal(event); err != nil {
return
}
route := fmt.Sprintf(
"%s.%s.%s.%d",
hostname,
event.ActionName,
event.Field,
event.GetUnitId())
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err = b.ch.PublishWithContext(ctx,
exchange,
route,
false,
false,
amqp.Publishing{
ContentType: "application/json",
Body: bytes,
}); err != nil {
fmt.Printf("Can't publish msg: %v", err)
}
}
func (b *broker) Start() {
var err error
b.conn, err = amqp.Dial(b.uri)
if err != nil {
panic(fmt.Sprintf("Can't connect to AMQP: %v", err))
}
b.ch, err = b.conn.Channel()
if err != nil {
panic(fmt.Sprintf("Can't create AMQP channel: %v", err))
}
fmt.Printf("AMQP client connected to %v, sending data to server..", b.uri)
}
func (b *broker) Stop() {
if b.conn != nil {
b.conn.Close()
}
if b.ch != nil {
b.ch.Close()
}
}
+35
View File
@@ -0,0 +1,35 @@
package ui
import (
"cli-mon/yabl"
"testing"
"time"
)
func TestNewAmqp(t *testing.T) {
t.Run("Amqp Test", func(t *testing.T) {
client := NewAmqp("amqp://user:password@localhost:5672/")
if client == nil {
t.Errorf("Can't create client")
}
client.Start()
now := time.Now()
event := &yabl.Event{
Field: yabl.FBoardReady,
ActionName: yabl.PPuErrors,
Object: nil,
//unit: 3,
Updated: &now,
Value: yabl.ERROR,
}
client.Consume(event)
client.Stop()
})
}