-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathhelper.go
More file actions
56 lines (47 loc) · 1.47 KB
/
helper.go
File metadata and controls
56 lines (47 loc) · 1.47 KB
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
package pubsub
import (
"cloud.google.com/go/pubsub/v2"
"context"
"github.com/xgodev/boost/bootstrap/function"
"github.com/xgodev/boost/wrapper/log"
"sync"
)
// Helper assists in creating event handlers for Pub/Sub with multiple topics.
type Helper[T any] struct {
handler function.Handler[T]
options *Options
client *pubsub.Client
}
// NewHelperWithOptions returns a new Helper with custom options.
func NewHelperWithOptions[T any](client *pubsub.Client, handler function.Handler[T], options *Options) *Helper[T] {
return &Helper[T]{
handler: handler,
options: options,
client: client,
}
}
// NewHelper returns a new Helper with default options.
func NewHelper[T any](client *pubsub.Client, handler function.Handler[T]) *Helper[T] {
opt, err := DefaultOptions()
if err != nil {
log.Fatal(err.Error())
}
return NewHelperWithOptions(client, handler, opt)
}
// Start subscribes to the topics and processes messages concurrently.
func (h *Helper[T]) Start() {
logger := log.WithTypeOf(*h)
var wg sync.WaitGroup
for _, subscription := range h.options.Subscriptions {
wg.Go(func() {
subscriber := NewSubscriber[T](h.client, h.handler, subscription, h.options)
if err := subscriber.Subscribe(context.Background()); err != nil {
logger.Errorf("Failed to subscribe to subscription %s: %v", subscription, err)
} else {
logger.Infof("Successfully subscribed to subscription %s", subscription)
}
})
}
// Wait for all subscriptions to complete
wg.Wait()
}