-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathnksh.go
59 lines (46 loc) · 1.21 KB
/
nksh.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
package nksh
import (
"context"
"os"
"os/signal"
"syscall"
"github.com/denkhaus/nksh/shared"
"github.com/juju/errors"
"github.com/sirupsen/logrus"
"golang.org/x/sync/errgroup"
)
var (
log logrus.FieldLogger = logrus.New().WithField("package", "nksh")
)
func Startup(kafkaHost, zookeeperHost string, funcs ...shared.DispatcherFunc) error {
kServers, err := LookupClusterHosts(kafkaHost, 9092)
if err != nil {
return errors.Annotate(err, "LookupClusterHosts [kafka]")
}
zServers, err := LookupClusterHosts(zookeeperHost, 2181)
if err != nil {
return errors.Annotate(err, "LookupClusterHosts [zookeeper]")
}
ctx, cancel := context.WithCancel(context.Background())
grp, ctx := errgroup.WithContext(ctx)
log.Infof("startup with kafka hosts %v", kServers)
log.Infof("startup with zookeeper hosts %v", zServers)
for _, fn := range funcs {
grp.Go(fn(ctx, kServers, zServers))
}
waiter := make(chan os.Signal, 1)
signal.Notify(waiter, syscall.SIGINT, syscall.SIGTERM)
select {
case <-waiter:
case <-ctx.Done():
}
cancel()
if err := grp.Wait(); err != nil {
return errors.Annotate(err, "Wait")
}
log.Info("dispatcher finished")
return nil
}
func SetLogger(logger logrus.FieldLogger) {
log = logger
}