Merge pull request #12 from galal-hussein/add_ingress_kubeconfig

Add ingress for clusters and kubeconfigs
This commit is contained in:
Hussein Galal
2023-01-20 02:30:43 +02:00
committed by GitHub
21 changed files with 612 additions and 73 deletions
+2
View File
@@ -25,6 +25,8 @@ spec:
type: integer
token:
type: string
ingressClassName:
type: string
scope: Namespaced
names:
plural: clusters
+28
View File
@@ -0,0 +1,28 @@
apiVersion: apps/v1
kind: Deployment
metadata:
name: k3k-controller
namespace: k3k-system
spec:
replicas: 1
selector:
matchLabels:
k8s-app: k3k-controller
template:
metadata:
labels:
k8s-app: k3k-controller
name: k3k-controller
spec:
containers:
- image: husseingalal/k3k:dev
imagePullPolicy: Always
name: k3k-controller
ports:
- containerPort: 8080
name: https
protocol: TCP
restartPolicy: Always
schedulerName: default-scheduler
serviceAccount: k3k-serviceaccount
serviceAccountName: k3k-serviceaccount
+4
View File
@@ -0,0 +1,4 @@
apiVersion: v1
kind: Namespace
metadata:
name: k3k-system
+12
View File
@@ -0,0 +1,12 @@
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
name: k3k-admin
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: ClusterRole
name: cluster-admin
subjects:
- kind: ServiceAccount
name: k3k-serviceaccount
namespace: k3k-system
+5
View File
@@ -0,0 +1,5 @@
apiVersion: v1
kind: ServiceAccount
metadata:
name: k3k-serviceaccount
namespace: k3k-system
@@ -1,10 +1,11 @@
apiVersion: k3k.io/v1alpha1
kind: Cluster
metadata:
name: example-cluster
name: multiple-servers
namespace: default
spec:
servers: 2
agents: 3
token: test
version: v1.26.0-k3s2
ingressClassName: traefik
+11
View File
@@ -0,0 +1,11 @@
apiVersion: k3k.io/v1alpha1
kind: Cluster
metadata:
name: single-server
namespace: default
spec:
servers: 1
agents: 3
token: test
version: v1.26.0-k3s2
ingressClassName: traefik
+8 -5
View File
@@ -3,9 +3,9 @@ module github.com/galal-hussein/k3k
go 1.19
require (
k8s.io/api v0.26.0
k8s.io/apimachinery v0.26.0
k8s.io/client-go v0.26.0
k8s.io/api v0.26.1
k8s.io/apimachinery v0.26.1
k8s.io/client-go v0.26.1
k8s.io/klog v1.0.0
)
@@ -34,6 +34,7 @@ require (
github.com/prometheus/client_model v0.3.0 // indirect
github.com/prometheus/common v0.37.0 // indirect
github.com/prometheus/procfs v0.8.0 // indirect
github.com/sirupsen/logrus v1.8.1 // indirect
github.com/spf13/pflag v1.0.5 // indirect
golang.org/x/oauth2 v0.0.0-20220223155221-ee480838109b // indirect
golang.org/x/sys v0.3.0 // indirect
@@ -44,7 +45,7 @@ require (
google.golang.org/protobuf v1.28.1 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
k8s.io/apiextensions-apiserver v0.26.0 // indirect
k8s.io/component-base v0.26.0 // indirect
k8s.io/component-base v0.26.1 // indirect
k8s.io/kube-openapi v0.0.0-20221012153701-172d655c2280 // indirect
sigs.k8s.io/yaml v1.3.0 // indirect
)
@@ -52,14 +53,16 @@ require (
require (
github.com/go-logr/logr v1.2.3 // indirect
github.com/gogo/protobuf v1.3.2 // indirect
github.com/google/gofuzz v1.1.0 // indirect
github.com/google/gofuzz v1.2.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/rancher/dynamiclistener v0.3.5
golang.org/x/net v0.3.1-0.20221206200815-1e63c2f08a10 // indirect
golang.org/x/text v0.5.0 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
k8s.io/apiserver v0.26.1
k8s.io/klog/v2 v2.80.1
k8s.io/utils v0.0.0-20221128185143-99ec85e7a448 // indirect
sigs.k8s.io/controller-runtime v0.14.1
+17 -10
View File
@@ -142,8 +142,8 @@ github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/
github.com/google/go-cmp v0.5.9 h1:O2Tfq5qg4qc4AmwVlvv0oLiVAGB7enBSJ2x2DqQFi38=
github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/gofuzz v1.1.0 h1:Hsa8mG0dQ46ij8Sl2AYJDUv1oA9/d6Vk+3LG99Oe02g=
github.com/google/gofuzz v1.1.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/gofuzz v1.2.0 h1:xRy4A+RhZaiKjJ1bPfwQ8sedCA+YS2YcCHW6ec7JMi0=
github.com/google/gofuzz v1.2.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/martian v2.1.0+incompatible/go.mod h1:9I4somxYTbIHy5NJKHRl3wXiIaQGbYVAs8BPL6v8lEs=
github.com/google/martian/v3 v3.0.0/go.mod h1:y5Zk1BBys9G+gd6Jrk0W3cC1+ELVxBWuIGO+w/tUAp0=
github.com/google/pprof v0.0.0-20181206194817-3ea8567a2e57/go.mod h1:zfwlbNMJ+OItoe0UupaVj+oy1omPYYDuagoSzA8v9mc=
@@ -241,10 +241,14 @@ github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1
github.com/prometheus/procfs v0.7.3/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA=
github.com/prometheus/procfs v0.8.0 h1:ODq8ZFEaYeCaZOJlZZdJA2AbQR98dSHSM1KW/You5mo=
github.com/prometheus/procfs v0.8.0/go.mod h1:z7EfXMXOkbkqb9IINtpCn86r/to3BnA0uaxHdg830/4=
github.com/rancher/dynamiclistener v0.3.5 h1:5TaIHvkDGmZKvc96Huur16zfTKOiLhDtK4S+WV0JA6A=
github.com/rancher/dynamiclistener v0.3.5/go.mod h1:dW/YF6/m2+uEyJ5VtEcd9THxda599HP6N9dSXk81+k0=
github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4=
github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo=
github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE=
github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88=
github.com/sirupsen/logrus v1.8.1 h1:dJKuHgqk1NNQlqoA6BTlM1Wf9DOH3NBjQyu0h9+AZZE=
github.com/sirupsen/logrus v1.8.1/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
github.com/spf13/pflag v1.0.5/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
github.com/stoewer/go-strcase v1.2.0/go.mod h1:IBiWB2sKIp3wVVQ3Y035++gc+knqhUQag1KpM8ahLw8=
@@ -370,6 +374,7 @@ golang.org/x/sys v0.0.0-20190606165138-5da285871e9c/go.mod h1:h1NjWce9XRLGQEsW7w
golang.org/x/sys v0.0.0-20190624142023-c5567b49c5d0/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190726091711-fc99dfbffb4e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191026070338-33540a1f6037/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191204072324-ce4227a45e2e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20191228213918-04cbcbbfeed8/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200106162015-b016eb3dc98e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
@@ -573,16 +578,18 @@ honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWh
honnef.co/go/tools v0.0.1-2019.2.3/go.mod h1:a3bituU0lyd329TUQxRnasdCoJDkEUEAqEt0JzvZhAg=
honnef.co/go/tools v0.0.1-2020.1.3/go.mod h1:X/FiERA/W4tHapMX5mGpAtMSVEeEUOyHaw9vFzvIQ3k=
honnef.co/go/tools v0.0.1-2020.1.4/go.mod h1:X/FiERA/W4tHapMX5mGpAtMSVEeEUOyHaw9vFzvIQ3k=
k8s.io/api v0.26.0 h1:IpPlZnxBpV1xl7TGk/X6lFtpgjgntCg8PJ+qrPHAC7I=
k8s.io/api v0.26.0/go.mod h1:k6HDTaIFC8yn1i6pSClSqIwLABIcLV9l5Q4EcngKnQg=
k8s.io/api v0.26.1 h1:f+SWYiPd/GsiWwVRz+NbFyCgvv75Pk9NK6dlkZgpCRQ=
k8s.io/api v0.26.1/go.mod h1:xd/GBNgR0f707+ATNyPmQ1oyKSgndzXij81FzWGsejg=
k8s.io/apiextensions-apiserver v0.26.0 h1:Gy93Xo1eg2ZIkNX/8vy5xviVSxwQulsnUdQ00nEdpDo=
k8s.io/apiextensions-apiserver v0.26.0/go.mod h1:7ez0LTiyW5nq3vADtK6C3kMESxadD51Bh6uz3JOlqWQ=
k8s.io/apimachinery v0.26.0 h1:1feANjElT7MvPqp0JT6F3Ss6TWDwmcjLypwoPpEf7zg=
k8s.io/apimachinery v0.26.0/go.mod h1:tnPmbONNJ7ByJNz9+n9kMjNP8ON+1qoAIIC70lztu74=
k8s.io/client-go v0.26.0 h1:lT1D3OfO+wIi9UFolCrifbjUUgu7CpLca0AD8ghRLI8=
k8s.io/client-go v0.26.0/go.mod h1:I2Sh57A79EQsDmn7F7ASpmru1cceh3ocVT9KlX2jEZg=
k8s.io/component-base v0.26.0 h1:0IkChOCohtDHttmKuz+EP3j3+qKmV55rM9gIFTXA7Vs=
k8s.io/component-base v0.26.0/go.mod h1:lqHwlfV1/haa14F/Z5Zizk5QmzaVf23nQzCwVOQpfC8=
k8s.io/apimachinery v0.26.1 h1:8EZ/eGJL+hY/MYCNwhmDzVqq2lPl3N3Bo8rvweJwXUQ=
k8s.io/apimachinery v0.26.1/go.mod h1:tnPmbONNJ7ByJNz9+n9kMjNP8ON+1qoAIIC70lztu74=
k8s.io/apiserver v0.26.1 h1:6vmnAqCDO194SVCPU3MU8NcDgSqsUA62tBUSWrFXhsc=
k8s.io/apiserver v0.26.1/go.mod h1:wr75z634Cv+sifswE9HlAo5FQ7UoUauIICRlOE+5dCg=
k8s.io/client-go v0.26.1 h1:87CXzYJnAMGaa/IDDfRdhTzxk/wzGZ+/HUQpqgVSZXU=
k8s.io/client-go v0.26.1/go.mod h1:IWNSglg+rQ3OcvDkhY6+QLeasV4OYHDjdqeWkDQZwGE=
k8s.io/component-base v0.26.1 h1:4ahudpeQXHZL5kko+iDHqLj/FSGAEUnSVO0EBbgDd+4=
k8s.io/component-base v0.26.1/go.mod h1:VHrLR0b58oC035w6YQiBSbtsf0ThuSwXP+p5dD/kAWU=
k8s.io/klog v1.0.0 h1:Pt+yjF5aB1xDSVbau4VsWe+dQNzA0qv1LlXdC2dF6Q8=
k8s.io/klog v1.0.0/go.mod h1:4Bi6QPql/J/LkTDqv7R/cd3hPo4k2DG6Ptcz060Ez5I=
k8s.io/klog/v2 v2.80.1 h1:atnLQ121W371wYYFawwYx1aEY2eUfs4l3J72wtgAwV4=
+2 -2
View File
@@ -6,7 +6,7 @@ import (
"flag"
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
"github.com/galal-hussein/k3k/pkg/controller"
"github.com/galal-hussein/k3k/pkg/controller/cluster"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"k8s.io/client-go/tools/clientcmd"
@@ -41,7 +41,7 @@ func main() {
if err != nil {
klog.Fatalf("Failed to create new controller runtime manager: %v", err)
}
if err := controller.Add(mgr); err != nil {
if err := cluster.Add(mgr); err != nil {
klog.Fatalf("Failed to add the new controller: %v", err)
}
+6 -5
View File
@@ -15,11 +15,12 @@ type Cluster struct {
}
type ClusterSpec struct {
Name string `json:"name"`
Version string `json:"version"`
Servers *int32 `json:"servers"`
Agents *int32 `json:"agents"`
Token string `json:"token"`
Name string `json:"name"`
Version string `json:"version"`
Servers *int32 `json:"servers"`
Agents *int32 `json:"agents"`
Token string `json:"token"`
IngressClassName string `json:"ingressClassName"`
}
// +k8s:deepcopy-gen:interfaces=k8s.io/apimachinery/pkg/runtime.Object
@@ -1,4 +1,4 @@
package controller
package agent
import (
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
@@ -8,8 +8,8 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
func agent(cluster *v1alpha1.Cluster) *apps.Deployment {
image := getImage(cluster)
func Agent(cluster *v1alpha1.Cluster) *apps.Deployment {
image := util.K3SImage(cluster)
name := "k3k-agent"
@@ -1,11 +1,13 @@
package controller
package cluster
import (
"context"
"fmt"
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
"github.com/galal-hussein/k3k/pkg/controller/config"
"github.com/galal-hussein/k3k/pkg/controller/cluster/agent"
"github.com/galal-hussein/k3k/pkg/controller/cluster/config"
"github.com/galal-hussein/k3k/pkg/controller/cluster/server"
"github.com/galal-hussein/k3k/pkg/controller/util"
v1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -24,7 +26,7 @@ const (
ClusterController = "k3k-cluster-controller"
)
type K3KReconciler struct {
type ClusterReconciler struct {
Client client.Client
Scheme *runtime.Scheme
}
@@ -32,7 +34,7 @@ type K3KReconciler struct {
// Add adds a new controller to the manager
func Add(mgr manager.Manager) error {
// initialize a new Reconciler
reconciler := K3KReconciler{
reconciler := ClusterReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
}
@@ -44,54 +46,48 @@ func Add(mgr manager.Manager) error {
MaxConcurrentReconciles: 1,
})
if err != nil {
return err
}
if err := controller.Watch(&source.Kind{Type: &v1alpha1.Cluster{}},
&handler.EnqueueRequestForObject{}); err != nil {
return err
}
return err
return nil
}
func (r *K3KReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) {
func (r *ClusterReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) {
cluster := v1alpha1.Cluster{}
if err := r.Client.Get(ctx, req.NamespacedName, &cluster); err != nil {
return reconcile.Result{}, r.wrapErr(fmt.Sprintf("failed to get cluster %s", req.NamespacedName), err)
return reconcile.Result{}, util.WrapErr(fmt.Sprintf("failed to get cluster %s", req.NamespacedName), err)
}
klog.Infof("%v", !cluster.DeletionTimestamp.IsZero())
if !cluster.DeletionTimestamp.IsZero() {
if err := r.handleDeletion(ctx, &cluster); err != nil {
return reconcile.Result{}, r.wrapErr(fmt.Sprintf("failed to delete cluster %s", req.NamespacedName), err)
return reconcile.Result{}, util.WrapErr(fmt.Sprintf("failed to delete cluster %s", req.NamespacedName), err)
}
}
// we create a namespace for each new cluster
ns := &v1.Namespace{}
if err := r.Client.Get(ctx, client.ObjectKey{Name: util.ClusterNamespace(&cluster)}, ns); err != nil {
if apierrors.IsNotFound(err) {
klog.Infof("creating new cluster")
return reconcile.Result{}, r.createCluster(ctx, &cluster)
} else {
if !apierrors.IsNotFound(err) {
return reconcile.Result{},
r.wrapErr(fmt.Sprintf("failed to get cluster namespace %s", util.ClusterNamespace(&cluster)), err)
util.WrapErr(fmt.Sprintf("failed to get cluster namespace %s", util.ClusterNamespace(&cluster)), err)
}
}
return reconcile.Result{}, nil
klog.Infof("enqueue cluster [%s]", cluster.Name)
return reconcile.Result{}, r.createCluster(ctx, &cluster)
}
// handleDeletion will delete the k3k cluster from kubernetes
func (r *K3KReconciler) handleDeletion(ctx context.Context, cluster *v1alpha1.Cluster) error {
func (r *ClusterReconciler) handleDeletion(ctx context.Context, cluster *v1alpha1.Cluster) error {
return nil
}
func (r *K3KReconciler) wrapErr(errString string, err error) error {
klog.Errorf("%s: %v", errString, err)
return err
}
func (r *K3KReconciler) createCluster(ctx context.Context, cluster *v1alpha1.Cluster) error {
func (r *ClusterReconciler) createCluster(ctx context.Context, cluster *v1alpha1.Cluster) error {
// create a new namespace for the cluster
namespace := &v1.Namespace{
ObjectMeta: metav1.ObjectMeta{
@@ -99,13 +95,18 @@ func (r *K3KReconciler) createCluster(ctx context.Context, cluster *v1alpha1.Clu
},
}
if err := r.Client.Create(ctx, namespace); err != nil {
return r.wrapErr("failed to create ns", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create ns", err)
}
}
// create cluster service
clusterService := service(cluster)
clusterService := server.Service(cluster)
if err := r.Client.Create(ctx, &clusterService); err != nil {
return r.wrapErr("failed to create cluster service", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create cluster service", err)
}
}
service := v1.Service{}
@@ -114,43 +115,75 @@ func (r *K3KReconciler) createCluster(ctx context.Context, cluster *v1alpha1.Clu
Namespace: util.ClusterNamespace(cluster),
Name: "k3k-server-service"},
&service); err != nil {
return r.wrapErr("failed to get cluster service", err)
return util.WrapErr("failed to get cluster service", err)
}
// create init node config
initServerConfigMap := config.ServerConfig(cluster, true, service.Spec.ClusterIP)
if err := r.Client.Create(ctx, &initServerConfigMap); err != nil {
return r.wrapErr("failed to create init configmap", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create init configmap", err)
}
}
// create servers configuration
serverConfigMap := config.ServerConfig(cluster, false, service.Spec.ClusterIP)
if err := r.Client.Create(ctx, &serverConfigMap); err != nil {
return r.wrapErr("failed to create configmap", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create configmap", err)
}
}
// create deployment for the init server
// the init deployment must have only 1 replica
initNodeDeployment := server(cluster, true)
initNodeDeployment := server.Server(cluster, true)
if err := r.Client.Create(ctx, initNodeDeployment); err != nil {
return r.wrapErr("failed to create init node deployment", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create init node deployment", err)
}
}
// create deployment for the rest of the servers
serverNodesDeployment := server(cluster, false)
serverNodesDeployment := server.Server(cluster, false)
if err := r.Client.Create(ctx, serverNodesDeployment); err != nil {
return r.wrapErr("failed to create server nodes deployment", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create server nodes deployment", err)
}
}
agentsConfigMap := config.AgentConfig(cluster, service.Spec.ClusterIP)
if err := r.Client.Create(ctx, &agentsConfigMap); err != nil {
return r.wrapErr("failed to create agent config", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create agent config", err)
}
}
agentsDeployment := agent(cluster)
agentsDeployment := agent.Agent(cluster)
if err := r.Client.Create(ctx, agentsDeployment); err != nil {
return r.wrapErr("failed to create agent deployment", err)
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create agent deployment", err)
}
}
// create ingress with random port for the server
serverIngress, err := server.Ingress(ctx, cluster, r.Client)
if err != nil {
return util.WrapErr("failed to create ingress object", err)
}
if err := r.Client.Create(ctx, serverIngress); err != nil {
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create server ingress", err)
}
}
kubeconfigSecret, err := server.GenerateNewKubeConfig(ctx, cluster, service.Spec.ClusterIP)
if err != nil {
return util.WrapErr("failed to generate new kubeconfig", err)
}
if err := r.Client.Create(ctx, kubeconfigSecret); err != nil {
if !apierrors.IsAlreadyExists(err) {
return util.WrapErr("failed to create kubeconfig secret", err)
}
}
return nil
}
+102
View File
@@ -0,0 +1,102 @@
package server
import (
"context"
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
"github.com/galal-hussein/k3k/pkg/controller/util"
v1 "k8s.io/api/core/v1"
networkingv1 "k8s.io/api/networking/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
)
var (
pathType = networkingv1.PathTypePrefix
wildcardDNS = ".sslip.io"
)
func Ingress(ctx context.Context, cluster *v1alpha1.Cluster, client client.Client) (*networkingv1.Ingress, error) {
addresses, err := addresses(ctx, client)
if err != nil {
return nil, err
}
ingressRules := ingressRules(cluster, addresses)
return &networkingv1.Ingress{
TypeMeta: metav1.TypeMeta{
Kind: "Ingress",
APIVersion: "networking.k8s.io/v1",
},
ObjectMeta: metav1.ObjectMeta{
Name: cluster.Name + "-server-ingress",
Namespace: util.ClusterNamespace(cluster),
},
Spec: networkingv1.IngressSpec{
IngressClassName: &cluster.Spec.IngressClassName,
Rules: ingressRules,
},
}, nil
}
// return all the nodes external addresses, if not found then return internal addresses
func addresses(ctx context.Context, client client.Client) ([]string, error) {
addresses := []string{}
nodeList := v1.NodeList{}
if err := client.List(ctx, &nodeList); err != nil {
return nil, err
}
for _, node := range nodeList.Items {
addresses = append(addresses, GetNodeAddress(&node))
}
return addresses, nil
}
func GetNodeAddress(node *v1.Node) string {
externalIP := ""
internalIP := ""
for _, ip := range node.Status.Addresses {
if ip.Type == "ExternalIP" && ip.Address != "" {
externalIP = ip.Address
break
} else if ip.Type == "InternalIP" && ip.Address != "" {
internalIP = ip.Address
}
}
if externalIP != "" {
return externalIP
}
return internalIP
}
func ingressRules(cluster *v1alpha1.Cluster, addresses []string) []networkingv1.IngressRule {
ingressRules := []networkingv1.IngressRule{}
for _, address := range addresses {
rule := networkingv1.IngressRule{
Host: cluster.Name + "." + address + wildcardDNS,
IngressRuleValue: networkingv1.IngressRuleValue{
HTTP: &networkingv1.HTTPIngressRuleValue{
Paths: []networkingv1.HTTPIngressPath{
{
Path: "/",
PathType: &pathType,
Backend: networkingv1.IngressBackend{
Service: &networkingv1.IngressServiceBackend{
Name: "k3k-server-service",
Port: networkingv1.ServiceBackendPort{
Number: 6443,
},
},
},
},
},
},
},
}
ingressRules = append(ingressRules, rule)
}
return ingressRules
}
+247
View File
@@ -0,0 +1,247 @@
package server
import (
"context"
"crypto"
"crypto/tls"
"crypto/x509"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
"github.com/galal-hussein/k3k/pkg/controller/util"
certutil "github.com/rancher/dynamiclistener/cert"
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apiserver/pkg/authentication/user"
"k8s.io/client-go/tools/clientcmd"
clientcmdapi "k8s.io/client-go/tools/clientcmd/api"
"k8s.io/client-go/util/retry"
)
const (
adminCommonName = "system:admin"
port = 6443
)
type controlRuntimeBootstrap struct {
ServerCA content
ServerCAKey content
ClientCA content
ClientCAKey content
}
type content struct {
Timestamp string
Content string
}
// GenerateNewKubeConfig generates the kubeconfig for the server:
// 1- use the server token to get the bootstrap data from k3s
// 2- generate client admin cert/key
// 3- use the ca cert from the bootstrap data & admin cert/key to write a new kubeconfig
// 4- save the new kubeconfig as a secret
func GenerateNewKubeConfig(ctx context.Context, cluster *v1alpha1.Cluster, ip string) (*v1.Secret, error) {
token := cluster.Spec.Token
bootstrap := &controlRuntimeBootstrap{}
err := retry.OnError(retry.DefaultBackoff, func(err error) bool {
return true
}, func() error {
var err error
bootstrap, err = requestBootstrap(token, ip)
return err
})
if err != nil {
return nil, err
}
if err := decodeBootstrap(bootstrap); err != nil {
return nil, err
}
adminCert, adminKey, err := createClientCertKey(
adminCommonName,
[]string{user.SystemPrivilegedGroup},
nil,
[]x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth},
bootstrap.ClientCA.Content,
bootstrap.ClientCAKey.Content)
if err != nil {
return nil, err
}
url := fmt.Sprintf("https://%s:%d", ip, port)
kubeconfigData, err := kubeconfig(url, []byte(bootstrap.ServerCA.Content), adminCert, adminKey)
if err != nil {
return nil, err
}
return &v1.Secret{
TypeMeta: metav1.TypeMeta{
Kind: "Secret",
APIVersion: "v1",
},
ObjectMeta: metav1.ObjectMeta{
Name: cluster.Name + "-kubeconfig",
Namespace: util.ClusterNamespace(cluster),
},
Data: map[string][]byte{
"kubeconfig.yaml": kubeconfigData,
},
}, nil
}
func requestBootstrap(token, serverIP string) (*controlRuntimeBootstrap, error) {
url := "https://" + serverIP + ":6443/v1-k3s/server-bootstrap"
client := http.Client{
Transport: &http.Transport{
TLSClientConfig: &tls.Config{
InsecureSkipVerify: true,
},
},
Timeout: 5 * time.Second,
}
req, err := http.NewRequest("GET", url, nil)
if err != nil {
return nil, err
}
req.Header.Add("Authorization", "Basic "+basicAuth("server", token))
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
runtimeBootstrap := controlRuntimeBootstrap{}
if err != nil {
return nil, err
}
if err := json.Unmarshal(body, &runtimeBootstrap); err != nil {
return nil, err
}
return &runtimeBootstrap, nil
}
func createClientCertKey(
commonName string,
organization []string,
altNames *certutil.AltNames,
extKeyUsage []x509.ExtKeyUsage,
caCert,
caKey string) ([]byte, []byte, error) {
caKeyPEM, err := certutil.ParsePrivateKeyPEM([]byte(caKey))
if err != nil {
return nil, nil, err
}
caCertPEM, err := certutil.ParseCertsPEM([]byte(caCert))
if err != nil {
return nil, nil, err
}
keyBytes, err := generateKey()
if err != nil {
return nil, nil, err
}
key, err := certutil.ParsePrivateKeyPEM(keyBytes)
if err != nil {
return nil, nil, err
}
cfg := certutil.Config{
CommonName: commonName,
Organization: organization,
Usages: extKeyUsage,
}
if altNames != nil {
cfg.AltNames = *altNames
}
cert, err := certutil.NewSignedCert(cfg, key.(crypto.Signer), caCertPEM[0], caKeyPEM.(crypto.Signer))
if err != nil {
return nil, nil, err
}
return append(certutil.EncodeCertPEM(cert), certutil.EncodeCertPEM(caCertPEM[0])...), keyBytes, nil
}
func generateKey() (data []byte, err error) {
generatedData, err := certutil.MakeEllipticPrivateKeyPEM()
if err != nil {
return nil, fmt.Errorf("error generating key: %v", err)
}
return generatedData, nil
}
func kubeconfig(url string, serverCA, clientCert, clientKey []byte) ([]byte, error) {
config := clientcmdapi.NewConfig()
cluster := clientcmdapi.NewCluster()
cluster.CertificateAuthorityData = serverCA
cluster.Server = url
authInfo := clientcmdapi.NewAuthInfo()
authInfo.ClientCertificateData = clientCert
authInfo.ClientKeyData = clientKey
context := clientcmdapi.NewContext()
context.AuthInfo = "default"
context.Cluster = "default"
config.Clusters["default"] = cluster
config.AuthInfos["default"] = authInfo
config.Contexts["default"] = context
config.CurrentContext = "default"
kubeconfig, err := clientcmd.Write(*config)
if err != nil {
return nil, err
}
return kubeconfig, nil
}
func basicAuth(username, password string) string {
auth := username + ":" + password
return base64.StdEncoding.EncodeToString([]byte(auth))
}
func decodeBootstrap(bootstrap *controlRuntimeBootstrap) error {
//client-ca
decoded, err := base64.StdEncoding.DecodeString(bootstrap.ClientCA.Content)
if err != nil {
return err
}
bootstrap.ClientCA.Content = string(decoded)
//client-ca-key
decoded, err = base64.StdEncoding.DecodeString(bootstrap.ClientCAKey.Content)
if err != nil {
return err
}
bootstrap.ClientCAKey.Content = string(decoded)
//server-ca
decoded, err = base64.StdEncoding.DecodeString(bootstrap.ServerCA.Content)
if err != nil {
return err
}
bootstrap.ServerCA.Content = string(decoded)
//server-ca-key
decoded, err = base64.StdEncoding.DecodeString(bootstrap.ServerCAKey.Content)
if err != nil {
return err
}
bootstrap.ServerCAKey.Content = string(decoded)
return nil
}
@@ -1,4 +1,4 @@
package controller
package server
import (
"strconv"
@@ -10,9 +10,9 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
func server(cluster *v1alpha1.Cluster, init bool) *apps.Deployment {
func Server(cluster *v1alpha1.Cluster, init bool) *apps.Deployment {
var replicas int32
image := getImage(cluster)
image := util.K3SImage(cluster)
name := "k3k-server"
if init {
@@ -55,10 +55,6 @@ func server(cluster *v1alpha1.Cluster, init bool) *apps.Deployment {
}
}
func getImage(cluster *v1alpha1.Cluster) string {
return "rancher/k3s:" + cluster.Spec.Version
}
func serverPodSpec(image, name string) v1.PodSpec {
privileged := true
return v1.PodSpec{
@@ -1,4 +1,4 @@
package controller
package server
import (
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
@@ -7,7 +7,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
func service(cluster *v1alpha1.Cluster) v1.Service {
func Service(cluster *v1alpha1.Cluster) v1.Service {
return v1.Service{
TypeMeta: metav1.TypeMeta{
Kind: "Service",
@@ -0,0 +1,74 @@
package ingressupdate
import (
"context"
"github.com/galal-hussein/k3k/pkg/controller/util"
v1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/manager"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
"sigs.k8s.io/controller-runtime/pkg/source"
)
const (
IngressUpdateController = "ingress-update-controller"
)
type IngressUpdateReconciler struct {
Client client.Client
Scheme *runtime.Scheme
}
// Add adds a new controller to the manager
func Add(mgr manager.Manager) error {
// initialize a new Reconciler
reconciler := IngressUpdateReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
}
controller, err := controller.New(IngressUpdateController, mgr, controller.Options{
Reconciler: &reconciler,
MaxConcurrentReconciles: 1,
})
if err != nil {
return err
}
return controller.Watch(&source.Kind{Type: &v1.Node{}}, &handler.EnqueueRequestForObject{})
}
// Reconcile will update ingresses each time a node added/deleted/changes with the new addresses
func (r *IngressUpdateReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) {
nodeList := v1.NodeList{}
if err := r.Client.List(ctx, &nodeList); err != nil {
return reconcile.Result{}, util.WrapErr("failed to list nodes", err)
}
// if err := r.Client.Get(ctx, req.NamespacedName, &cluster); err != nil {
// return reconcile.Result{}, util.WrapErr(fmt.Sprintf("failed to get cluster %s", req.NamespacedName), err)
// }
// // we create a namespace for each new cluster
// ns := &v1.Namespace{}
// if err := r.Client.Get(ctx, client.ObjectKey{Name: util.ClusterNamespace(&cluster)}, ns); err != nil {
// if apierrors.IsNotFound(err) {
// klog.Infof("creating new cluster")
// return reconcile.Result{}, r.createCluster(ctx, &cluster)
// } else {
// return reconcile.Result{},
// util.WrapErr(fmt.Sprintf("failed to get cluster namespace %s", util.ClusterNamespace(&cluster)), err)
// }
// }
return reconcile.Result{}, nil
}
+14 -1
View File
@@ -1,11 +1,24 @@
package util
import "github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
import (
"github.com/galal-hussein/k3k/pkg/apis/k3k.io/v1alpha1"
"k8s.io/klog"
)
const (
NamespacePrefix = "k3k-"
K3SImageName = "rancher/k3s"
)
func ClusterNamespace(cluster *v1alpha1.Cluster) string {
return NamespacePrefix + cluster.Name
}
func K3SImage(cluster *v1alpha1.Cluster) string {
return K3SImageName + ":" + cluster.Spec.Version
}
func WrapErr(errString string, err error) error {
klog.Errorf("%s: %v", errString, err)
return err
}