-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy patheventmanager.go
78 lines (65 loc) · 1.49 KB
/
eventmanager.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
package eventmanager
import (
"encoding/json"
"fmt"
cloudevents "github.com/cloudevents/sdk-go/v2"
"github.com/google/uuid"
"github.com/nats-io/nats.go"
)
type EventHandler func(m *cloudevents.Event)
type EventManager struct {
Broker *nats.Conn
subscriptions map[string]*nats.Subscription
app string
}
func New(natsAddr string, app string) (*EventManager, error) {
nc, err := nats.Connect(natsAddr)
if err != nil {
return nil, err
}
em := &EventManager{
Broker: nc,
subscriptions: make(map[string]*nats.Subscription),
app: app,
}
return em, nil
}
func natsWrap(f EventHandler) nats.MsgHandler {
return func(m *nats.Msg) {
evt := cloudevents.NewEvent()
if err := json.Unmarshal(m.Data, &evt); err != nil {
return
}
f(&evt)
return
}
}
func (m *EventManager) Listen(subject string, cb EventHandler) error {
sub, err := m.Broker.Subscribe(subject, natsWrap(cb))
if err != nil {
return err
}
m.subscriptions[subject] = sub
return nil
}
func (m *EventManager) Publish(subject, t string, data interface{}) error {
evt := cloudevents.NewEvent()
u, err := uuid.NewUUID()
if err != nil {
return err
}
evt.SetID(u.String())
evt.SetSource(fmt.Sprintf("auttaja.io/%s", m.app))
evt.SetType(t)
if err := evt.SetData(cloudevents.ApplicationJSON, data); err != nil {
return err
}
bytes, err := json.Marshal(evt)
if err != nil {
return err
}
if err := m.Broker.Publish(subject, bytes); err != nil {
return err
}
return nil
}