/
main.go
74 lines (55 loc) · 1.46 KB
/
main.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
package main
import (
"context"
"fmt"
"log"
"os"
"os/signal"
"syscall"
"github.com/Shopify/sarama"
"github.com/syncromatics/kafql/internal/decoder"
"github.com/syncromatics/kafql/internal/kafka"
"github.com/syncromatics/kafql/internal/server"
v1 "github.com/syncromatics/proto-schema-registry/pkg/proto/schema/registry/v1"
"golang.org/x/sync/errgroup"
"google.golang.org/grpc"
)
func main() {
settings, err := getSettingsFromEnv()
if err != nil {
log.Fatal(err)
}
ctx, cancel := context.WithCancel(context.Background())
group, ctx := errgroup.WithContext(ctx)
con, err := grpc.Dial(settings.ProtoRegistry, grpc.WithInsecure())
if err != nil {
log.Fatal(err)
}
client := v1.NewRegistryAPIClient(con)
proto := decoder.NewProto(client)
conf := sarama.NewConfig()
conf.Version = sarama.MaxVersion
admin, err := sarama.NewClusterAdmin([]string{settings.KafkaBroker}, conf)
if err != nil {
log.Fatal(err)
}
consumer, err := sarama.NewConsumer([]string{settings.KafkaBroker}, conf)
if err != nil {
log.Fatal(err)
}
service := kafka.NewService(admin, consumer, proto)
server := server.NewServer(service)
group.Go(server.Start(ctx, "/", settings.Port))
eventChan := make(chan os.Signal)
signal.Notify(eventChan, syscall.SIGINT, syscall.SIGTERM)
fmt.Println("server started...")
select {
case <-eventChan:
case <-ctx.Done():
}
fmt.Println("server stopping...")
cancel()
if err := group.Wait(); err != nil {
log.Fatal(err)
}
}