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) } type config struct { PVE pve.Config `json:"pve"` } 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 cr, err := os.Open(cloudConfig) if err == nil { buf, _ := ioutil.ReadAll(cr) if len(buf) != 0 { var cfg 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.PVE) 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("approver"), 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", "approver") 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") */ }