-
Notifications
You must be signed in to change notification settings - Fork 35
Expand file tree
/
Copy pathmain.go
More file actions
110 lines (96 loc) · 2.71 KB
/
Copy pathmain.go
File metadata and controls
110 lines (96 loc) · 2.71 KB
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
package main
import (
"context"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
"github.com/apple/fdb-joshua/k8s/agent-scaler/internal/config"
"github.com/apple/fdb-joshua/k8s/agent-scaler/internal/joshua"
"github.com/apple/fdb-joshua/k8s/agent-scaler/internal/scaler"
)
func main() {
ctrl.SetLogger(zap.New())
log := ctrl.Log.WithName("agent-scaler")
cfg, err := config.Load()
if err != nil {
log.Error(err, "loading config")
os.Exit(1)
}
scheme := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(scheme); err != nil {
log.Error(err, "adding client-go scheme")
os.Exit(1)
}
if err := batchv1.AddToScheme(scheme); err != nil {
log.Error(err, "adding batchv1 scheme")
os.Exit(1)
}
if err := corev1.AddToScheme(scheme); err != nil {
log.Error(err, "adding corev1 scheme")
os.Exit(1)
}
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
Scheme: scheme,
Metrics: metricsserver.Options{
BindAddress: cfg.MetricsAddr,
},
HealthProbeBindAddress: cfg.HealthAddr,
LeaderElection: false,
})
if err != nil {
log.Error(err, "creating manager")
os.Exit(1)
}
counter, err := joshua.NewEnsembleCounter(cfg.FDBClusterFile, cfg.JoshuaNamespace)
if err != nil {
log.Error(err, "opening FDB connection")
os.Exit(1)
}
s := scaler.New(cfg, mgr.GetClient(), counter, log.WithName("scaler"))
if err := mgr.Add(s); err != nil {
log.Error(err, "adding scaler to manager")
os.Exit(1)
}
if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
log.Error(err, "adding healthz check")
os.Exit(1)
}
if err := mgr.AddReadyzCheck("readyz", func(_ *http.Request) error {
if !s.IsReady() {
return fmt.Errorf("not ready")
}
return nil
}); err != nil {
log.Error(err, "adding readyz check")
os.Exit(1)
}
log.Info("starting agent-scaler",
"agentName", cfg.AgentName,
"maxJobs", cfg.MaxJobs,
"batchSize", cfg.BatchSize,
"checkDelay", cfg.CheckDelay,
"namespace", cfg.Namespace,
)
// SIGKILL cannot be intercepted — the OS delivers it unconditionally.
// We handle SIGTERM (Kubernetes graceful stop) and SIGINT (Ctrl-C).
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
go func() {
<-ctx.Done()
log.Info("shutdown signal received, draining")
}()
if err := mgr.Start(ctx); err != nil {
log.Error(err, "manager exited")
os.Exit(1)
}
}