This repository has been archived by the owner on Jun 20, 2019. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathutil.go
85 lines (74 loc) · 1.85 KB
/
util.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
package main
import (
"crypto/tls"
"crypto/x509"
"fmt"
"io/ioutil"
"os"
"strings"
"gopkg.in/Shopify/sarama.v1"
)
func brokers(useTLS bool) []string {
var s string
if useTLS {
s = os.Getenv("TLS_KAFKA_BROKERS")
} else {
s = os.Getenv("KAFKA_BROKERS")
}
if s == "" {
s = "127.0.0.1:9092"
}
return strings.Split(s, ",")
}
func offsets(client sarama.Client, topic string, partition int32) (oldest int64, newest int64) {
oldest, err := client.GetOffset(topic, partition, sarama.OffsetOldest)
must(err)
newest, err = client.GetOffset(topic, partition, sarama.OffsetNewest)
must(err)
return oldest, newest
}
func tlsConfig() (useTLS bool, config *tls.Config, err error) {
// if SSL_CA_BUNDLE_PATH isn't set, just don't use TLS at all
caPath := os.Getenv("SSL_CA_BUNDLE_PATH")
if caPath == "" {
return
}
useTLS = true
config = new(tls.Config)
caCerts, err := ioutil.ReadFile(caPath)
if err != nil {
err = fmt.Errorf("error reading $SSL_CA_BUNDLE_PATH: %v", err)
return
}
config.RootCAs = x509.NewCertPool()
if !config.RootCAs.AppendCertsFromPEM(caCerts) {
err = fmt.Errorf("$SSL_CA_BUNDLE_PATH=%q was empty", caPath)
return
}
// if $SSL_CERT_PATH or $SSL_KEY_PATH aren't set, skip client cert
certPath, keyPath := os.Getenv("SSL_CRT_PATH"), os.Getenv("SSL_KEY_PATH")
if certPath == "" || keyPath == "" {
fmt.Println("SSL_CRT_PATH or SSL_KEY_PATH was empty!")
return
}
keypair, err := tls.LoadX509KeyPair(certPath, keyPath)
if err != nil {
err = fmt.Errorf("error reading $SSL_CRT_PATH/$SSL_KEY_PATH: %v", err)
return
}
config.Certificates = []tls.Certificate{keypair}
return
}
func getKafkaVersion(version string) (v sarama.KafkaVersion) {
switch version {
case "0.9.0.0":
v = sarama.V0_9_0_0
case "0.9.0.1":
v = sarama.V0_9_0_1
case "0.10.0.0":
v = sarama.V0_10_0_0
default:
v = sarama.V0_10_0_0
}
return
}