-
-
Notifications
You must be signed in to change notification settings - Fork 26
/
pubsub.go
64 lines (49 loc) · 1.3 KB
/
pubsub.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
// Package pubsub is the GCP PubSub implementation of an event stream.
package pubsub
import (
"context"
"fmt"
"cloud.google.com/go/pubsub"
"github.com/italolelis/outboxer"
)
const (
// TopicNameOption is the topic name option.
TopicNameOption = "topic_name"
// OrderingKeyOption is the ordering key option.
OrderingKeyOption = "ordering_key"
)
// Pubsub is the wrapper for the Pubsub library.
type Pubsub struct {
client *pubsub.Client
}
type options struct {
topicName string
orderingKey string
}
// New creates a new instance of Kinesis.
func New(c *pubsub.Client) *Pubsub {
return &Pubsub{client: c}
}
// Send sends the message to the event stream.
func (p *Pubsub) Send(ctx context.Context, evt *outboxer.OutboxMessage) error {
opts := p.parseOptions(evt.Options)
topic := p.client.Topic(opts.topicName)
res := topic.Publish(ctx, &pubsub.Message{
Data: evt.Payload,
OrderingKey: opts.orderingKey,
})
if _, err := res.Get(ctx); err != nil {
return fmt.Errorf("error when getting generated id: %v", err)
}
return nil
}
func (p *Pubsub) parseOptions(opts outboxer.DynamicValues) *options {
opt := options{}
if data, ok := opts[TopicNameOption]; ok {
opt.topicName = data.(string)
}
if data, ok := opts[OrderingKeyOption]; ok {
opt.orderingKey = data.(string)
}
return &opt
}