cloud-provider-pve/cmd/pve-csr-approver-manager/main.go
2020-11-01 01:17:54 +01:00

220 lines
5.8 KiB
Go

package main
import (
"encoding/json"
"errors"
"flag"
"io/ioutil"
"math/rand"
"os"
"time"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
typedcertificatesv1beta1 "k8s.io/client-go/kubernetes/typed/certificates/v1beta1"
"k8s.io/klog/klogr"
"k8s.io/klog/v2"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/controller"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"git.giftfish.de/ston1th/cloud-controller-manager-pve/pkg/approver"
pve "git.giftfish.de/ston1th/pve-go"
)
var (
scheme = runtime.NewScheme()
setupLog = ctrl.Log.WithName("setup")
)
func init() {
_ = clientgoscheme.AddToScheme(scheme)
}
func main() {
rand.Seed(time.Now().UnixNano())
klog.InitFlags(nil)
var (
leaderElect bool
healthAddr string
cloudConfig string
csrConcurrency int
syncPeriod time.Duration
)
flag.BoolVar(&leaderElect, "leader-elect", true,
"Enable leader election for controller manager. "+
"Enabling this will ensure there is only one active controller manager.")
flag.StringVar(&cloudConfig,
"cloud-config",
"",
"The path to the cloud provider configuration file. Empty string for no configuration file.",
)
flag.StringVar(&healthAddr,
"health-addr",
":9449",
"The address the health endpoint binds to.",
)
flag.IntVar(&csrConcurrency,
"csr-concurrency",
1,
"Number of CSRs to process simultaneously",
)
flag.DurationVar(&syncPeriod,
"sync-period",
10*time.Minute,
"The minimum interval at which watched resources are reconciled (e.g. 15m)",
)
flag.Parse()
ctrl.SetLogger(klogr.New())
rcfg := ctrl.GetConfigOrDie()
mgr, err := ctrl.NewManager(rcfg, ctrl.Options{
Scheme: scheme,
LeaderElection: leaderElect,
LeaderElectionID: "pve-csr-approver-manager",
HealthProbeBindAddress: healthAddr,
MetricsBindAddress: "0",
SyncPeriod: &syncPeriod,
})
if err != nil {
setupLog.Error(err, "unable to start manager")
os.Exit(1)
}
var opts []pve.ClientOption
config, err := os.Open(cloudConfig)
if err == nil {
buf, _ := ioutil.ReadAll(config)
if len(buf) != 0 {
var cfg pve.Config
err = json.Unmarshal(buf, &cfg)
if err != nil {
setupLog.Error(err, "failed to read config file")
os.Exit(1)
}
opts, err = pve.ClientOptionsFromConfig(cfg)
if err != nil {
setupLog.Error(err, "failed to create client options from file")
os.Exit(1)
}
}
} else {
opts, err = pve.ClientOptionsFromEnv()
if err != nil {
setupLog.Error(err, "failed to create client options from env")
os.Exit(1)
}
}
if len(opts) == 0 {
setupLog.Error(errors.New("missing pve config"), "")
os.Exit(1)
}
pveClient := pve.NewClient(opts...)
if err = (&approver.CSRReconciler{
Client: mgr.GetClient(),
Log: ctrl.Log.WithName("controllers").WithName("CSR"),
Scheme: mgr.GetScheme(),
PVEClient: pveClient,
CSRClient: typedcertificatesv1beta1.NewForConfigOrDie(rcfg),
}).SetupWithManager(mgr, controller.Options{MaxConcurrentReconciles: csrConcurrency}); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "CSR")
os.Exit(1)
}
if err := mgr.AddReadyzCheck("ping", healthz.Ping); err != nil {
setupLog.Error(err, "unable to create ready check")
os.Exit(1)
}
if err := mgr.AddHealthzCheck("ping", healthz.Ping); err != nil {
setupLog.Error(err, "unable to create health check")
os.Exit(1)
}
setupLog.Info("starting manager")
if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
setupLog.Error(err, "problem running manager")
os.Exit(1)
}
/*
cfg, err := rest.InClusterConfig()
if err != nil {
log.Error(err, "could not configure kubernetes client")
os.Exit(1)
}
cl, err := kubernetes.NewForConfig(cfg)
if err != nil {
log.Error(err, "could not create kubernetes client")
os.Exit(1)
}
id, err := os.Hostname()
if err != nil {
log.Error(err, "could not get hostname")
os.Exit(1)
}
lock := &resourcelock.LeaseLock{
LeaseMeta: metav1.ObjectMeta{
Name: csrLock,
Namespace: namespace,
},
Client: cl.CoordinationV1(),
LockConfig: resourcelock.ResourceLockConfig{
Identity: id,
},
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{
Lock: lock,
ReleaseOnCancel: true,
LeaseDuration: time.Duration(leaseDuration) * time.Second,
RenewDeadline: time.Duration(renewDeadline) * time.Second,
RetryPeriod: time.Duration(retryPeriod) * time.Second,
Callbacks: leaderelection.LeaderCallbacks{
OnStartedLeading: func(ctx context.Context) {
log.Info("Leader election successful")
factory := informers.NewSharedInformerFactory(cl, time.Second*30)
certInformerV1beta1 := factory.Certificates().V1beta1().CertificateSigningRequests().Informer()
fV1beta1 := func(obj interface{}) {
if req, ok := obj.(*certificatesv1beta1.CertificateSigningRequest); ok {
if err := approver.ApproveV1beta1(log, pveClient, cl.CertificatesV1beta1().CertificateSigningRequests(), req); err != nil {
log.Error(err, "csr approval failed", "name", req.ObjectMeta.Name)
return
}
}
}
certInformerV1beta1.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
fV1beta1(obj)
},
UpdateFunc: func(_, obj interface{}) {
fV1beta1(obj)
},
})
stop := make(chan struct{})
factory.Start(stop)
<-ctx.Done()
close(stop)
},
OnStoppedLeading: func() {
cancel()
log.Info("Leader election lost")
},
}})
log.Info("csr approver started")
sigs := make(chan os.Signal)
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
<-sigs
cancel()
log.Info("csr approver stopped")
*/
}