@@ -2,31 +2,59 @@ package processor
22
33import (
44 "context"
5+ "encoding/json"
6+ "fmt"
57 "time"
68
79 "go.uber.org/zap"
810 "google.golang.org/protobuf/proto"
911
10- "github.com/bucketeer-io/bucketeer/pkg/notification/sender"
12+ "github.com/bucketeer-io/bucketeer/pkg/notification/sender/notifier "
1113 "github.com/bucketeer-io/bucketeer/pkg/pubsub/puller"
1214 "github.com/bucketeer-io/bucketeer/pkg/pubsub/puller/codes"
1315 "github.com/bucketeer-io/bucketeer/pkg/subscriber"
1416 domainevent "github.com/bucketeer-io/bucketeer/proto/event/domain"
1517 domaineventproto "github.com/bucketeer-io/bucketeer/proto/event/domain"
18+ notificationproto "github.com/bucketeer-io/bucketeer/proto/notification"
19+ senderproto "github.com/bucketeer-io/bucketeer/proto/notification/sender"
1620)
1721
22+ type DemoOrganizationCreationNotifierConfig struct {
23+ WebURL string `json:"webURL"`
24+ SlackWebhookURL string `json:"slackWebhookURL"`
25+ }
26+
1827type demoOrganizationCreationNotifier struct {
19- sender sender.Sender
20- logger * zap.Logger
28+ slackNotifier notifier.Notifier
29+ demoOrganizationCreationNotifierConfig DemoOrganizationCreationNotifierConfig
30+ logger * zap.Logger
2131}
2232
2333func NewDemoOrganizationCreationNotifier (
24- sender sender. Sender ,
34+ config interface {} ,
2535 logger * zap.Logger ,
2636) subscriber.PubSubProcessor {
37+ jsonConfigMap , ok := config .(map [string ]interface {})
38+ if ! ok {
39+ logger .Error ("demoOrganizationCreationNotifier: invalid config type, expected map[string]interface{}" )
40+ return nil
41+ }
42+ configBytes , err := json .Marshal (jsonConfigMap )
43+ if err != nil {
44+ logger .Error ("demoOrganizationCreationNotifier: failed to marshal config" , zap .Error (err ))
45+ return nil
46+ }
47+ var notifierConfig DemoOrganizationCreationNotifierConfig
48+ if err := json .Unmarshal (configBytes , & notifierConfig ); err != nil {
49+ logger .Error ("demoOrganizationCreationNotifier: failed to unmarshal config" , zap .Error (err ))
50+ return nil
51+ }
52+ slackNotifier := notifier .NewSlackNotifier (notifierConfig .WebURL )
53+
2754 return & demoOrganizationCreationNotifier {
28- sender : sender ,
29- logger : logger ,
55+ slackNotifier : slackNotifier ,
56+ demoOrganizationCreationNotifierConfig : notifierConfig ,
57+ logger : logger ,
3058 }
3159}
3260
@@ -55,6 +83,11 @@ func (d demoOrganizationCreationNotifier) handleMessage(msg *puller.Message) {
5583 }
5684 domainEvent , err := d .unmarshalMessage (msg )
5785 if err != nil {
86+ d .logger .Error ("Failed to unmarshal message" ,
87+ zap .Error (err ),
88+ zap .String ("msgID" , msg .ID ),
89+ zap .String ("attributes" , fmt .Sprintf ("%+v" , msg .Attributes )),
90+ )
5891 subscriberHandledCounter .WithLabelValues (subscriberDemoOrganizationEvent , codes .BadMessage .String ()).Inc ()
5992 msg .Ack ()
6093 return
@@ -67,11 +100,42 @@ func (d demoOrganizationCreationNotifier) handleMessage(msg *puller.Message) {
67100 return
68101 }
69102
70- err = d .SendSlackNotifier (ctx , domainEvent )
103+ var organizationCreatedEvent domaineventproto.OrganizationCreatedEvent
104+ if err := domainEvent .Data .UnmarshalTo (& organizationCreatedEvent ); err != nil {
105+ d .logger .Error ("Failed to unmarshal OrganizationCreatedEvent" ,
106+ zap .String ("event id" , domainEvent .Id ),
107+ zap .Error (err ),
108+ )
109+ subscriberHandledCounter .WithLabelValues (
110+ subscriberDemoOrganizationEvent ,
111+ codes .NonRepeatableError .String (),
112+ ).Inc ()
113+ msg .Ack ()
114+ return
115+ }
116+
117+ recipient := & notificationproto.Recipient {
118+ Type : notificationproto .Recipient_SlackChannel ,
119+ Language : notificationproto .Recipient_ENGLISH ,
120+ SlackChannelRecipient : & notificationproto.SlackChannelRecipient {
121+ WebhookUrl : d .demoOrganizationCreationNotifierConfig .SlackWebhookURL ,
122+ },
123+ }
124+ fmt .Printf ("?%+v\n " , recipient )
125+ err = d .slackNotifier .Notify (ctx , & senderproto.Notification {
126+ Type : senderproto .Notification_DemoOrganizationCreation ,
127+ DemoOrganizationCreationNotification : & senderproto.DemoOrganizationCreationNotification {
128+ OwnerEmail : organizationCreatedEvent .OwnerEmail ,
129+ OrganizationId : organizationCreatedEvent .Id ,
130+ OrganizationName : organizationCreatedEvent .Name ,
131+ },
132+ }, recipient , recipient .Language )
71133 if err != nil {
72- d .logger .Error ("Failed to send Slack notification" ,
134+ d .logger .Error ("Failed to send notification" ,
73135 zap .Error (err ),
74- zap .String ("eventID" , domainEvent .Id ),
136+ zap .String ("event id" , domainEvent .Id ),
137+ zap .String ("webhookURL" , d .demoOrganizationCreationNotifierConfig .SlackWebhookURL ),
138+ zap .String ("organizationId" , organizationCreatedEvent .Id ),
75139 )
76140 subscriberHandledCounter .WithLabelValues (
77141 subscriberDemoOrganizationEvent ,
@@ -80,16 +144,14 @@ func (d demoOrganizationCreationNotifier) handleMessage(msg *puller.Message) {
80144 msg .Ack ()
81145 return
82146 }
147+ fmt .Printf ("?2" )
148+ subscriberHandledCounter .WithLabelValues (
149+ subscriberDemoOrganizationEvent ,
150+ codes .OK .String (),
151+ ).Inc ()
83152 msg .Ack ()
84153}
85154
86- func (d demoOrganizationCreationNotifier ) SendSlackNotifier (
87- ctx context.Context ,
88- event * domainevent.Event ,
89- ) error {
90-
91- }
92-
93155func (d demoOrganizationCreationNotifier ) unmarshalMessage (msg * puller.Message ) (* domainevent.Event , error ) {
94156 event := & domaineventproto.Event {}
95157 err := proto .Unmarshal (msg .Data , event )
0 commit comments