-
Notifications
You must be signed in to change notification settings - Fork 50
/
publisher.go
100 lines (81 loc) · 3.21 KB
/
publisher.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
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
package mdlsub
import (
"context"
"fmt"
"strconv"
"github.com/justtrackio/gosoline/pkg/cfg"
"github.com/justtrackio/gosoline/pkg/log"
"github.com/justtrackio/gosoline/pkg/mdl"
"github.com/justtrackio/gosoline/pkg/stream"
)
const (
AttributeModelId = "modelId"
AttributeType = "type"
AttributeVersion = "version"
ConfigKeyMdlSubPublishers = "mdlsub.publishers"
TypeCreate = "create"
TypeUpdate = "update"
TypeDelete = "delete"
)
type PublisherSettings struct {
mdl.ModelId
Producer string `cfg:"producer" validate:"required_without=OutputType"`
OutputType string `cfg:"output_type" validate:"required_without=Producer"`
Shared bool `cfg:"shared"`
}
//go:generate mockery --name Publisher
type Publisher interface {
PublishBatch(ctx context.Context, typ string, version int, values []interface{}, customAttributes ...map[string]string) error
Publish(ctx context.Context, typ string, version int, value interface{}, customAttributes ...map[string]string) error
}
type publisher struct {
logger log.Logger
producer stream.Producer
settings *PublisherSettings
}
func NewPublisher(ctx context.Context, config cfg.Config, logger log.Logger, name string) (*publisher, error) {
settings := readPublisherSetting(config, name)
return NewPublisherWithSettings(ctx, config, logger, settings)
}
func NewPublisherWithSettings(ctx context.Context, config cfg.Config, logger log.Logger, settings *PublisherSettings) (*publisher, error) {
var err error
var producer stream.Producer
if producer, err = stream.NewProducer(ctx, config, logger, settings.Producer); err != nil {
return nil, fmt.Errorf("can not create producer %s: %w", settings.Producer, err)
}
return NewPublisherWithInterfaces(logger, producer, settings), nil
}
func NewPublisherWithInterfaces(logger log.Logger, producer stream.Producer, settings *PublisherSettings) *publisher {
return &publisher{
logger: logger,
producer: producer,
settings: settings,
}
}
func (p *publisher) PublishBatch(ctx context.Context, typ string, version int, values []interface{}, customAttributes ...map[string]string) error {
attributes := []map[string]string{
CreateMessageAttributes(p.settings.ModelId, typ, version),
}
attributes = append(attributes, customAttributes...)
if err := p.producer.Write(ctx, values, attributes...); err != nil {
return fmt.Errorf("can not publish %s with publisher %s: %w", p.settings.ModelId.String(), p.settings.Name, err)
}
return nil
}
func (p *publisher) Publish(ctx context.Context, typ string, version int, value interface{}, customAttributes ...map[string]string) error {
attributes := []map[string]string{
CreateMessageAttributes(p.settings.ModelId, typ, version),
}
attributes = append(attributes, customAttributes...)
if err := p.producer.WriteOne(ctx, value, attributes...); err != nil {
return fmt.Errorf("can not publish %s with publisher %s: %w", p.settings.ModelId.String(), p.settings.Name, err)
}
return nil
}
func CreateMessageAttributes(modelId mdl.ModelId, typ string, version int) map[string]string {
return map[string]string{
AttributeType: typ,
AttributeVersion: strconv.Itoa(version),
AttributeModelId: modelId.String(),
}
}