csr rewrite
This commit is contained in:
parent
211b89236e
commit
954b5c0d61
8 changed files with 430 additions and 280 deletions
|
|
@ -1,104 +1,220 @@
|
|||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"io/ioutil"
|
||||
"math/rand"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
certificatesv1beta1 "k8s.io/api/certificates/v1beta1"
|
||||
"k8s.io/client-go/informers"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
"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()
|
||||
log := klogr.New()
|
||||
opts, err := pve.ClientOptionsFromEnv()
|
||||
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,
|
||||
SyncPeriod: &syncPeriod,
|
||||
})
|
||||
if err != nil {
|
||||
log.Error(err, "could not create pve client")
|
||||
setupLog.Error(err, "unable to start manager")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
pveClient := pve.NewClient(opts...)
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
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
|
||||
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)
|
||||
}
|
||||
}
|
||||
certInformerV1beta1.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
fV1beta1(obj)
|
||||
},
|
||||
UpdateFunc: func(_, obj interface{}) {
|
||||
fV1beta1(obj)
|
||||
},
|
||||
})
|
||||
stop := make(chan struct{})
|
||||
factory.Start(stop)
|
||||
if len(opts) == 0 {
|
||||
setupLog.Error(errors.New("missing pve config"), "")
|
||||
os.Exit(1)
|
||||
}
|
||||
pveClient := pve.NewClient(opts...)
|
||||
|
||||
//watchListV1 := cache.NewListWatchFromClient(
|
||||
// cl.CertificatesV1Client.RESTClient(),
|
||||
// "certificatesigningrequests",
|
||||
// v1.NamespaceAll,
|
||||
// fields.Everything(),
|
||||
//)
|
||||
//fV1 := func(obj interface{}) {
|
||||
// if req, ok := obj.(*certificates.CertificateSigningRequest); ok {
|
||||
// if err := approver.ApproveV1(pveClient, cl.CertificatesV1Client.CertificateSigningRequests(), req); err != nil {
|
||||
// log.Error(err, "csr approval failed", "name", req.ObjectMeta.Name)
|
||||
// return
|
||||
// }
|
||||
// log.Info("csr approval successful", "name", req.ObjectMeta.Name)
|
||||
// }
|
||||
//}
|
||||
//_, controllerV1 := cache.NewInformer(
|
||||
// watchListV1,
|
||||
// &certificates.CertificateSigningRequest{},
|
||||
// time.Second*30,
|
||||
// cache.ResourceEventHandlerFuncs{
|
||||
// AddFunc: func(obj interface{}) {
|
||||
// fV1(obj)
|
||||
// },
|
||||
// UpdateFunc: func(_, obj interface{}) {
|
||||
// fV1(obj)
|
||||
// },
|
||||
// },
|
||||
//)
|
||||
//stopV1 := make(chan struct{})
|
||||
//go controllerV1.Run(stopV1)
|
||||
log.Info("csr approver started")
|
||||
if err = (&approver.CSRReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Log: ctrl.Log.WithName("controllers").WithName("CSR"),
|
||||
Scheme: mgr.GetScheme(),
|
||||
Recorder: mgr.GetEventRecorderFor("pvecluster-controller"),
|
||||
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)
|
||||
}
|
||||
|
||||
sigs := make(chan os.Signal)
|
||||
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
|
||||
<-sigs
|
||||
close(stop)
|
||||
log.Info("csr approver stopped")
|
||||
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")
|
||||
*/
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue