-
Notifications
You must be signed in to change notification settings - Fork 13
/
Copy pathconfig.go
139 lines (121 loc) · 4.44 KB
/
config.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
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
package goneli
import (
"fmt"
"os"
"path/filepath"
"regexp"
"time"
validation "github.com/go-ozzo/ozzo-validation"
"github.com/obsidiandynamics/libstdgo/scribe"
)
// Duration is a convenience for deriving a pointer from a given Duration argument.
func Duration(d time.Duration) *time.Duration {
return &d
}
func defaultDuration(d **time.Duration, def time.Duration) {
if *d == nil {
*d = &def
}
}
// KafkaConfigMap represents the Kafka key-value configuration.
type KafkaConfigMap map[string]interface{}
// Config encapsulates configuration for Neli.
type Config struct {
KafkaConfig KafkaConfigMap
LeaderTopic string
LeaderGroupID string
KafkaConsumerProvider KafkaConsumerProvider
KafkaProducerProvider KafkaProducerProvider
Scribe scribe.Scribe
Name string
PollDuration *time.Duration
MinPollInterval *time.Duration
HeartbeatTimeout *time.Duration
}
// Validate the Config, returning an error if invalid.
func (c Config) Validate() error {
validName := matchValidKafkaChars()
return validation.ValidateStruct(&c,
validation.Field(&c.KafkaConfig, validation.NotNil),
validation.Field(&c.LeaderTopic, validation.Required, validation.Match(validName)),
validation.Field(&c.LeaderGroupID, validation.Required, validation.Match(validName)),
validation.Field(&c.KafkaConsumerProvider, validation.NotNil),
validation.Field(&c.KafkaProducerProvider, validation.NotNil),
validation.Field(&c.Scribe, validation.NotNil),
validation.Field(&c.Name, validation.Required, validation.Match(regexp.MustCompile("^[^%]*$"))),
validation.Field(&c.PollDuration, validation.Required, validation.Min(1*time.Millisecond)),
validation.Field(&c.MinPollInterval, validation.Required, validation.Min(1*time.Millisecond)),
validation.Field(&c.HeartbeatTimeout, validation.Required, validation.Min(1*time.Millisecond)),
)
}
// Obtains a textual representation of the configuration.
func (c Config) String() string {
return fmt.Sprint(
"Config[KafkaConfig=", c.KafkaConfig,
", LeaderTopic=", c.LeaderTopic,
", LeaderGroupID=", c.LeaderGroupID,
", KafkaConsumerProvider=", c.KafkaConsumerProvider,
", KafkaProducerProvider=", c.KafkaProducerProvider,
", Scribe=", c.Scribe,
", Name=", c.Name,
", PollDuration=", c.PollDuration,
", MinPollInterval=", c.MinPollInterval,
", HeartbeatTimeout=", c.HeartbeatTimeout, "]")
}
const (
// DefaultPollDuration is the default value of Config.PollDuration
DefaultPollDuration = 1 * time.Millisecond
// DefaultMinPollInterval is the default value of Config.MinPollInterval
DefaultMinPollInterval = 100 * time.Millisecond
// DefaultHeartbeatTimeout is the default value of Config.HeartbeatTimeout
DefaultHeartbeatTimeout = 5 * time.Second
)
// SetDefaults assigns the default values to optional fields.
func (c *Config) SetDefaults() {
if c.KafkaConfig == nil {
c.KafkaConfig = KafkaConfigMap{}
}
if _, ok := c.KafkaConfig["bootstrap.servers"]; !ok {
c.KafkaConfig["bootstrap.servers"] = "localhost:9092"
}
if c.LeaderGroupID == "" {
c.LeaderGroupID = Sanitise(filepath.Base(os.Args[0]))
}
if c.LeaderTopic == "" {
c.LeaderTopic = c.LeaderGroupID + ".neli"
}
if c.KafkaConsumerProvider == nil {
c.KafkaConsumerProvider = StandardKafkaConsumerProvider()
}
if c.KafkaProducerProvider == nil {
c.KafkaProducerProvider = StandardKafkaProducerProvider()
}
if c.Scribe == nil {
c.Scribe = scribe.New(scribe.StandardBinding())
}
if c.Name == "" {
c.Name = fmt.Sprintf("%s_%d_%d", Sanitise(getString("localhost", os.Hostname)), os.Getpid(), time.Now().Unix())
}
defaultDuration(&c.PollDuration, DefaultPollDuration)
defaultDuration(&c.MinPollInterval, DefaultMinPollInterval)
defaultDuration(&c.HeartbeatTimeout, DefaultHeartbeatTimeout)
}
const validKafkaNameChars = "a-zA-Z0-9\\._\\-"
// Obtains a regular expression matcher for the set of valid characters
// allowed in Kafka resource names.
func matchValidKafkaChars() *regexp.Regexp {
return regexp.MustCompile("^[" + validKafkaNameChars + "]*$")
}
// Sanitise cleans up the given name by replacing all characters in the pattern [^a-zA-Z0-9\\._\\-]
// with underscores.
func Sanitise(name string) string {
return string(regexp.MustCompile("[^"+validKafkaNameChars+"]").ReplaceAll([]byte(name), []byte("_")))
}
type stringGetter func() (string, error)
func getString(def string, stringGetter stringGetter) string {
str, err := stringGetter()
if err != nil {
return def
}
return str
}