Add network router and resource (#189)

This change adds two new resources, NetworkRouter and NetworkResource,
which enable clusters to expose Kubernetes services to Netbird.

The NetworkRouter is responsible for creating the network, group, setup
key and routing peer all of which are unique to the isntance. Along with
the deployment of the client in the cluster.

The NetworkResource exposes a service by linking to the specific router
it wants to expose to. This makes coupling between the resource and
network easy to understand.

Routers also set a DNS zone which is used to give names to resources
based on the name and namespace of the service being exposed.

Part of #172

Signed-off-by: Philip Laine <philip.laine@gmail.com>
This commit is contained in:
Philip Laine
2026-04-23 08:55:49 +02:00
committed by GitHub
parent af11e31b28
commit 6768a76c9c
38 changed files with 2973 additions and 55 deletions
@@ -0,0 +1,301 @@
package controller
import (
"context"
"fmt"
"strings"
"time"
"github.com/fluxcd/pkg/runtime/conditions"
"github.com/fluxcd/pkg/runtime/patch"
netbird "github.com/netbirdio/netbird/shared/management/client/rest"
"github.com/netbirdio/netbird/shared/management/http/api"
corev1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/builder"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
"github.com/netbirdio/kubernetes-operator/internal/netbirdutil"
)
type NetworkResourceReconciler struct {
client.Client
Netbird *netbird.Client
}
// +kubebuilder:rbac:groups=netbird.io,resources=networkresources,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=netbird.io,resources=networkresources/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=netbird.io,resources=networkresources/finalizers,verbs=update
// nolint:gocyclo
func (r *NetworkResourceReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
netResource := &nbv1alpha1.NetworkResource{}
err := r.Get(ctx, req.NamespacedName, netResource)
if err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
sp := patch.NewSerialPatcher(netResource, r.Client)
if !netResource.DeletionTimestamp.IsZero() {
return r.reconcileDelete(ctx, sp, netResource)
}
svc := &corev1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: netResource.Spec.ServiceRef.Name,
Namespace: netResource.Namespace,
},
}
err = r.Get(ctx, client.ObjectKeyFromObject(svc), svc)
if err != nil {
if kerrors.IsNotFound(err) {
conditions.MarkFalse(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.DependencyReason, "Referenced Service cannot be found.")
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
return ctrl.Result{}, err
}
if svc.Spec.Type != corev1.ServiceTypeClusterIP {
conditions.MarkFalse(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.DependencyReason, "Referenced Service is not of type ClusterIP.")
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
if svc.Spec.ClusterIP == "" || svc.Spec.ClusterIP == corev1.ClusterIPNone {
conditions.MarkFalse(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.DependencyReason, "Referenced Service does not have a ClusterIP set.")
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
netRouter := &nbv1alpha1.NetworkRouter{
ObjectMeta: metav1.ObjectMeta{
Name: netResource.Spec.NetworkRouterRef.Name,
Namespace: netResource.Spec.NetworkRouterRef.Namespace,
},
}
err = r.Get(ctx, client.ObjectKeyFromObject(netRouter), netRouter)
if err != nil {
if kerrors.IsNotFound(err) {
conditions.MarkFalse(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.DependencyReason, "Referenced NetworkRouter cannot be found.")
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
return ctrl.Result{}, err
}
if netRouter.Status.NetworkID == "" || netRouter.Status.RoutingPeerID == "" {
conditions.MarkFalse(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.DependencyReason, "Referenced NetworkRouter is not ready.")
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
groupIDs, err := netbirdutil.GetGroupIDs(ctx, r.Client, r.Netbird, netResource.Spec.Groups, netResource.Namespace)
if err != nil {
return ctrl.Result{}, err
}
controllerutil.AddFinalizer(netResource, nbv1alpha1.NetbirdFinalizer)
resourceID, err := func() (string, error) {
netReq := api.NetworkResourceRequest{
Name: svc.Name + "/" + svc.Namespace,
Address: svc.Spec.ClusterIP,
Enabled: true,
Groups: groupIDs,
}
if netResource.Status.ResourceID != "" {
netResp, err := r.Netbird.Networks.Resources(netRouter.Status.NetworkID).Update(ctx, netResource.Status.ResourceID, netReq)
if err != nil && !netbird.IsNotFound(err) {
return "", err
}
if err == nil {
return netResp.Id, nil
}
}
netResp, err := r.Netbird.Networks.Resources(netRouter.Status.NetworkID).Create(ctx, netReq)
if err != nil {
return "", err
}
return netResp.Id, nil
}()
if err != nil {
return ctrl.Result{}, err
}
netResource.Status.NetworkID = netRouter.Status.NetworkID
netResource.Status.ResourceID = resourceID
err = sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
// Create DNS records for resource.
zone, err := netbirdutil.GetDNSZoneByName(ctx, r.Netbird, netRouter.Spec.DNSZoneRef.Name)
if err != nil {
return ctrl.Result{}, err
}
// If zone has changed we need to delete the old records.
if netResource.Status.DNSZoneID != "" && netResource.Status.DNSZoneID != zone.Id {
err = r.Netbird.DNSZones.DeleteRecord(ctx, netResource.Status.DNSZoneID, netResource.Status.DNSRecordID)
if err != nil && !netbird.IsNotFound(err) {
return ctrl.Result{}, err
}
netResource.Status.DNSZoneID = ""
netResource.Status.DNSRecordID = ""
}
recordID, err := func() (string, error) {
dnsReq := api.DNSRecordRequest{
Content: svc.Spec.ClusterIP,
Name: strings.Join([]string{svc.Name, svc.Namespace, zone.Name}, "."),
Ttl: int(5 * time.Minute / time.Second),
Type: api.DNSRecordTypeA,
}
if netResource.Status.DNSZoneID != "" && netResource.Status.DNSRecordID != "" {
recordResp, err := r.Netbird.DNSZones.UpdateRecord(ctx, netResource.Status.DNSZoneID, netResource.Status.DNSRecordID, dnsReq)
if err != nil && !netbird.IsNotFound(err) {
return "", err
}
if err == nil {
return recordResp.Id, nil
}
}
recordResp, err := r.Netbird.DNSZones.CreateRecord(ctx, zone.Id, dnsReq)
if err != nil {
return "", err
}
return recordResp.Id, nil
}()
if err != nil {
return ctrl.Result{}, err
}
netResource.Status.DNSZoneID = zone.Id
netResource.Status.DNSRecordID = recordID
conditions.MarkTrue(netResource, nbv1alpha1.ReadyCondition, nbv1alpha1.ReconciledReason, "")
err = sp.Patch(ctx, netResource, patch.WithStatusObservedGeneration{})
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
func (r *NetworkResourceReconciler) reconcileDelete(ctx context.Context, sp *patch.SerialPatcher, netResource *nbv1alpha1.NetworkResource) (ctrl.Result, error) {
if netResource.Status.NetworkID != "" && netResource.Status.ResourceID != "" {
err := r.Netbird.Networks.Resources(netResource.Status.NetworkID).Delete(ctx, netResource.Status.ResourceID)
if err != nil && !netbird.IsNotFound(err) {
return ctrl.Result{}, err
}
}
if netResource.Status.DNSZoneID != "" && netResource.Status.DNSRecordID != "" {
err := r.Netbird.DNSZones.DeleteRecord(ctx, netResource.Status.DNSZoneID, netResource.Status.DNSRecordID)
if err != nil && !netbird.IsNotFound(err) {
return ctrl.Result{}, err
}
}
controllerutil.RemoveFinalizer(netResource, nbv1alpha1.NetbirdFinalizer)
err := sp.Patch(ctx, netResource)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
func (r *NetworkResourceReconciler) SetupWithManager(mgr ctrl.Manager) error {
err := mgr.GetFieldIndexer().IndexField(context.Background(), &nbv1alpha1.NetworkResource{}, ".spec.networkRouterRef", func(obj client.Object) []string {
netResource := obj.(*nbv1alpha1.NetworkResource)
ref := netResource.Spec.NetworkRouterRef
if ref.Name == "" {
return nil
}
if ref.Namespace == "" {
ref.Namespace = netResource.Namespace
}
return []string{fmt.Sprintf("%s/%s", ref.Name, ref.Namespace)}
})
if err != nil {
return err
}
err = mgr.GetFieldIndexer().IndexField(context.Background(), &nbv1alpha1.NetworkResource{}, ".spec.serviceRef", func(obj client.Object) []string {
netResource := obj.(*nbv1alpha1.NetworkResource)
ref := netResource.Spec.ServiceRef
if ref.Name == "" {
return nil
}
return []string{netResource.Spec.ServiceRef.Name}
})
if err != nil {
return err
}
return ctrl.NewControllerManagedBy(mgr).
For(&nbv1alpha1.NetworkResource{}).
Watches(
&nbv1alpha1.NetworkRouter{},
handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, obj client.Object) []reconcile.Request {
netResourceList := &nbv1alpha1.NetworkResourceList{}
err := r.List(ctx, netResourceList, client.MatchingFields{".spec.networkRouterRef": fmt.Sprintf("%s/%s", obj.GetName(), obj.GetNamespace())})
if err != nil {
return nil
}
requests := make([]reconcile.Request, len(netResourceList.Items))
for i, item := range netResourceList.Items {
requests[i] = reconcile.Request{
NamespacedName: types.NamespacedName{
Name: item.Name,
Namespace: item.Namespace,
},
}
}
return requests
}),
builder.WithPredicates(predicate.ResourceVersionChangedPredicate{}),
).
Watches(
&corev1.Service{},
handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, obj client.Object) []reconcile.Request {
netResourceList := &nbv1alpha1.NetworkResourceList{}
err := r.List(ctx, netResourceList, client.InNamespace(obj.GetNamespace()), client.MatchingFields{".spec.serviceRef": obj.GetName()})
if err != nil {
return nil
}
requests := make([]reconcile.Request, len(netResourceList.Items))
for i, item := range netResourceList.Items {
requests[i] = reconcile.Request{
NamespacedName: types.NamespacedName{
Name: item.Name,
Namespace: item.Namespace,
},
}
}
return requests
}),
builder.WithPredicates(predicate.ResourceVersionChangedPredicate{}),
).
Complete(r)
}
@@ -0,0 +1,150 @@
package controller
import (
"context"
"fmt"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
corev1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
"github.com/netbirdio/kubernetes-operator/internal/netbirdmock"
"github.com/netbirdio/netbird/shared/management/http/api"
)
var _ = Describe("NetworkResource Controller", func() {
Context("When reconciling a resource", func() {
ctx := context.Background()
var netResourceRec *NetworkResourceReconciler
var netRouterRec *NetworkRouterReconciler
var setupKeyRec *SetupKeyReconciler
var groupRec *GroupReconciler
nn := client.ObjectKey{
Name: "test-resource",
Namespace: "network-resource",
}
BeforeEach(func() {
nbClient := netbirdmock.Client()
netResourceRec = &NetworkResourceReconciler{
Client: k8sClient,
Netbird: nbClient,
}
netRouterRec = &NetworkRouterReconciler{
Client: k8sClient,
Netbird: nbClient,
ClientImage: "docker.io/netbirdio/netbird:latest",
ManagementURL: "https://netbird.io",
}
setupKeyRec = &SetupKeyReconciler{
Client: k8sClient,
Netbird: nbClient,
}
groupRec = &GroupReconciler{
Client: k8sClient,
Netbird: nbClient,
}
ns := &corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Namespace,
},
}
Expect(k8sClient.Create(ctx, ns)).To(Succeed())
})
AfterEach(func() {
ns := &corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Namespace,
},
}
err := k8sClient.Get(ctx, client.ObjectKeyFromObject(ns), ns)
if kerrors.IsNotFound(err) {
return
}
Expect(err).ToNot(HaveOccurred())
Expect(k8sClient.Delete(ctx, ns)).To(Succeed())
})
It("creates a network resource and DNS record", func() {
zoneReq := api.ZoneRequest{
Name: "cluster.local",
Domain: "cluster.local",
}
_, err := netRouterRec.Netbird.DNSZones.CreateZone(ctx, zoneReq)
Expect(err).ToNot(HaveOccurred())
// Create network router that we reference.
netRouter := &nbv1alpha1.NetworkRouter{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Name,
Namespace: nn.Namespace,
},
Spec: nbv1alpha1.NetworkRouterSpec{
DNSZoneRef: nbv1alpha1.DNSZoneReference{
Name: "cluster.local",
},
},
}
Expect(k8sClient.Create(ctx, netRouter)).To(Succeed())
for range 3 {
_, err := netRouterRec.Reconcile(ctx, reconcile.Request{NamespacedName: nn})
Expect(err).NotTo(HaveOccurred())
key := client.ObjectKey{Name: fmt.Sprintf("networkrouter-%s", netRouter.Name), Namespace: nn.Namespace}
_, err = groupRec.Reconcile(ctx, reconcile.Request{NamespacedName: key})
Expect(err).NotTo(HaveOccurred())
_, err = setupKeyRec.Reconcile(ctx, reconcile.Request{NamespacedName: key})
Expect(err).NotTo(HaveOccurred())
}
svc := &corev1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: nn.Namespace,
},
Spec: corev1.ServiceSpec{
Ports: []corev1.ServicePort{
{
Port: 8080,
},
},
},
}
Expect(k8sClient.Create(ctx, svc)).To(Succeed())
netResource := &nbv1alpha1.NetworkResource{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Name,
Namespace: nn.Namespace,
},
Spec: nbv1alpha1.NetworkResourceSpec{
NetworkRouterRef: nbv1alpha1.CrossNamespaceReference{
Name: netRouter.Name,
Namespace: netRouter.Namespace,
},
ServiceRef: corev1.LocalObjectReference{
Name: svc.Name,
},
},
}
Expect(k8sClient.Create(ctx, netResource)).To(Succeed())
_, err = netResourceRec.Reconcile(ctx, reconcile.Request{NamespacedName: nn})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, nn, netResource)
Expect(err).NotTo(HaveOccurred())
Expect(netResource.Status.NetworkID).NotTo(BeEmpty())
Expect(netResource.Status.ResourceID).NotTo(BeEmpty())
Expect(netResource.Status.DNSZoneID).NotTo(BeEmpty())
Expect(netResource.Status.DNSRecordID).NotTo(BeEmpty())
})
})
})
@@ -0,0 +1,322 @@
package controller
import (
"context"
"crypto/sha256"
"encoding/json"
"fmt"
"maps"
"time"
"github.com/fluxcd/pkg/runtime/conditions"
"github.com/fluxcd/pkg/runtime/patch"
"github.com/netbirdio/kubernetes-operator/internal/netbirdutil"
"github.com/netbirdio/kubernetes-operator/internal/ssautil"
netbird "github.com/netbirdio/netbird/shared/management/client/rest"
"github.com/netbirdio/netbird/shared/management/http/api"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/strategicpatch"
appsv1ac "k8s.io/client-go/applyconfigurations/apps/v1"
corev1ac "k8s.io/client-go/applyconfigurations/core/v1"
metav1ac "k8s.io/client-go/applyconfigurations/meta/v1"
"k8s.io/utils/ptr"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
nbv1alpha1ac "github.com/netbirdio/kubernetes-operator/pkg/applyconfigurations/api/v1alpha1"
)
type NetworkRouterReconciler struct {
client.Client
Netbird *netbird.Client
ManagementURL string
ClientImage string
}
// +kubebuilder:rbac:groups=netbird.io,resources=networkrouters,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=netbird.io,resources=networkrouters/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=netbird.io,resources=networkrouters/finalizers,verbs=update
// nolint:gocyclo
func (r *NetworkRouterReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
netRouter := &nbv1alpha1.NetworkRouter{}
err := r.Get(ctx, req.NamespacedName, netRouter)
if err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
sp := patch.NewSerialPatcher(netRouter, r.Client)
if !netRouter.DeletionTimestamp.IsZero() {
return r.reconcileDelete(ctx, sp, netRouter)
}
ownerRef, err := ssautil.OwnerReference(netRouter, r.Scheme())
if err != nil {
return ctrl.Result{}, err
}
// Ensure the DNS Zone exists.
_, err = netbirdutil.GetDNSZoneByName(ctx, r.Netbird, netRouter.Spec.DNSZoneRef.Name)
if err != nil {
return ctrl.Result{}, err
}
controllerutil.AddFinalizer(netRouter, nbv1alpha1.NetbirdFinalizer)
networkID, err := func() (string, error) {
networkReq := api.NetworkRequest{
Name: netRouter.Name,
}
if netRouter.Status.NetworkID != "" {
networkResp, err := r.Netbird.Networks.Update(ctx, netRouter.Status.NetworkID, networkReq)
if err != nil && !netbird.IsNotFound(err) {
return "", err
}
if err == nil {
return networkResp.Id, nil
}
}
networkResp, err := r.Netbird.Networks.Create(ctx, networkReq)
if err != nil {
return "", err
}
return networkResp.Id, nil
}()
if err != nil {
return ctrl.Result{}, err
}
netRouter.Status.NetworkID = networkID
err = sp.Patch(ctx, netRouter)
if err != nil {
return ctrl.Result{}, err
}
// Calculate unique suffix used for Netbird resources.
sum := sha256.Sum256([]byte(netRouter.UID))
uniqueSuffix := networkID + "-" + fmt.Sprintf("%x", sum[:4])[:8]
// Create the group used by the router to discover peers.
groupAC := nbv1alpha1ac.Group(fmt.Sprintf("networkrouter-%s", netRouter.Name), req.Namespace).
WithOwnerReferences(ownerRef).
WithSpec(
nbv1alpha1ac.GroupSpec().
WithName(fmt.Sprintf("networkrouter-%s", uniqueSuffix)),
)
err = r.Client.Apply(ctx, groupAC)
if err != nil {
return ctrl.Result{}, err
}
group := &nbv1alpha1.Group{
ObjectMeta: metav1.ObjectMeta{
Name: *groupAC.Name,
Namespace: *groupAC.Namespace,
},
}
err = r.Client.Get(ctx, client.ObjectKeyFromObject(group), group)
if err != nil {
return ctrl.Result{}, err
}
if group.Status.GroupID == "" {
return ctrl.Result{}, nil
}
// Create the setup key used by routing peers.
setupKeyAC := nbv1alpha1ac.SetupKey(fmt.Sprintf("networkrouter-%s", netRouter.Name), req.Namespace).
WithOwnerReferences(ownerRef).
WithSpec(
nbv1alpha1ac.SetupKeySpec().
WithName(fmt.Sprintf("networkrouter-%s", uniqueSuffix)).
WithEphemeral(true).
WithAutoGroups(nbv1alpha1ac.ResourceReference().WithID(group.Status.GroupID)),
)
err = r.Client.Apply(ctx, setupKeyAC)
if err != nil {
return ctrl.Result{}, err
}
setupKey := nbv1alpha1.SetupKey{
ObjectMeta: metav1.ObjectMeta{
Name: *setupKeyAC.Name,
Namespace: *setupKeyAC.Namespace,
},
}
err = r.Get(ctx, client.ObjectKeyFromObject(&setupKey), &setupKey)
if err != nil {
return ctrl.Result{}, err
}
if setupKey.Status.SetupKeyID == "" {
return ctrl.Result{}, nil
}
// Create the routing peer in netbird.
routingPeerID, err := func() (string, error) {
routerReq := api.NetworkRouterRequest{
Enabled: true,
Masquerade: true,
Metric: 9999,
PeerGroups: ptr.To([]string{group.Status.GroupID}),
}
if netRouter.Status.RoutingPeerID != "" {
resp, err := r.Netbird.Networks.Routers(networkID).Update(ctx, netRouter.Status.RoutingPeerID, routerReq)
if err != nil && !netbird.IsNotFound(err) {
return "", err
}
if err == nil {
return resp.Id, nil
}
}
resp, err := r.Netbird.Networks.Routers(networkID).Create(ctx, routerReq)
if err != nil {
return "", err
}
return resp.Id, nil
}()
if err != nil {
return ctrl.Result{}, err
}
netRouter.Status.RoutingPeerID = routingPeerID
err = sp.Patch(ctx, netRouter, patch.WithStatusObservedGeneration{})
if err != nil {
return ctrl.Result{}, err
}
// Create the deployment.
selectorLabels := map[string]string{
"app.kubernetes.io/name": "networkrouter",
"app.kubernetes.io/instance": req.Name,
}
podTemplateSpecAC := corev1ac.PodTemplateSpec().
WithLabels(selectorLabels).
WithSpec(corev1ac.PodSpec().
WithContainers(corev1ac.Container().
WithName("netbird").
WithImage(r.ClientImage).
WithEnv(
corev1ac.EnvVar().
WithName("NB_SETUP_KEY").
WithValueFrom(corev1ac.EnvVarSource().
WithSecretKeyRef(corev1ac.SecretKeySelector().
WithName(setupKey.SecretName()).
WithKey(SetupKeySecretKey),
),
),
corev1ac.EnvVar().
WithName("NB_MANAGEMENT_URL").
WithValue(r.ManagementURL),
corev1ac.EnvVar().
WithName("NB_LOG_LEVEL").
WithValue("info"),
).
WithStartupProbe(corev1ac.Probe().WithExec(corev1ac.ExecAction().WithCommand("netbird", "status", "--check", "startup"))).
WithReadinessProbe(corev1ac.Probe().WithExec(corev1ac.ExecAction().WithCommand("netbird", "status", "--check", "ready"))).
WithSecurityContext(corev1ac.SecurityContext().
WithCapabilities(corev1ac.Capabilities().
WithAdd("NET_ADMIN").
WithAdd("SYS_RESOURCE").
WithAdd("SYS_ADMIN"),
).
WithPrivileged(true),
),
),
)
depLabels := map[string]string{}
depAnnotations := map[string]string{}
replicas := int32(3)
if netRouter.Spec.WorkloadOverride != nil {
if netRouter.Spec.WorkloadOverride.Labels != nil {
depLabels = netRouter.Spec.WorkloadOverride.Labels
}
if netRouter.Spec.WorkloadOverride.Annotations != nil {
depAnnotations = netRouter.Spec.WorkloadOverride.Annotations
}
if netRouter.Spec.WorkloadOverride.Replicas != nil {
replicas = *netRouter.Spec.WorkloadOverride.Replicas
}
if netRouter.Spec.WorkloadOverride.PodTemplate != nil {
baseJSON, err := json.Marshal(&podTemplateSpecAC)
if err != nil {
return ctrl.Result{}, err
}
overrideJSON, err := json.Marshal(netRouter.Spec.WorkloadOverride.PodTemplate)
if err != nil {
return ctrl.Result{}, err
}
mergedJSON, err := strategicpatch.StrategicMergePatch(baseJSON, overrideJSON, corev1.PodTemplateSpec{})
if err != nil {
return ctrl.Result{}, err
}
err = json.Unmarshal(mergedJSON, &podTemplateSpecAC)
if err != nil {
return ctrl.Result{}, err
}
}
}
maps.Copy(depLabels, selectorLabels)
depAC := appsv1ac.Deployment(fmt.Sprintf("networkrouter-%s", req.Name), req.Namespace).
WithOwnerReferences(ownerRef).
WithLabels(depLabels).
WithAnnotations(depAnnotations).
WithSpec(appsv1ac.DeploymentSpec().WithReplicas(replicas).WithSelector(metav1ac.LabelSelector().WithMatchLabels(selectorLabels)).WithTemplate(podTemplateSpecAC))
err = r.Client.Apply(ctx, depAC)
if err != nil {
return ctrl.Result{}, err
}
dep := &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: *depAC.Name,
Namespace: *depAC.Namespace,
},
}
err = r.Client.Get(ctx, client.ObjectKeyFromObject(dep), dep)
if err != nil {
return ctrl.Result{}, err
}
if dep.Status.ReadyReplicas != dep.Status.Replicas {
return ctrl.Result{}, nil
}
conditions.MarkTrue(netRouter, nbv1alpha1.ReadyCondition, nbv1alpha1.ReconciledReason, "")
err = sp.Patch(ctx, netRouter, patch.WithStatusObservedGeneration{})
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{RequeueAfter: 15 * time.Minute}, nil
}
func (r *NetworkRouterReconciler) reconcileDelete(ctx context.Context, sp *patch.SerialPatcher, netRouter *nbv1alpha1.NetworkRouter) (ctrl.Result, error) {
if netRouter.Status.RoutingPeerID != "" {
err := r.Netbird.Networks.Routers(netRouter.Status.NetworkID).Delete(ctx, netRouter.Status.RoutingPeerID)
if err != nil && !netbird.IsNotFound(err) {
return ctrl.Result{}, err
}
}
if netRouter.Status.NetworkID != "" {
err := r.Netbird.Networks.Delete(ctx, netRouter.Status.NetworkID)
if err != nil && !netbird.IsNotFound(err) {
return ctrl.Result{}, err
}
}
controllerutil.RemoveFinalizer(netRouter, nbv1alpha1.NetbirdFinalizer)
err := sp.Patch(ctx, netRouter)
if err != nil {
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
func (r *NetworkRouterReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&nbv1alpha1.NetworkRouter{}).
Owns(&nbv1alpha1.Group{}).
Owns(&nbv1alpha1.SetupKey{}).
Owns(&appsv1.Deployment{}).
Complete(r)
}
@@ -0,0 +1,147 @@
package controller
import (
"context"
"fmt"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
"github.com/netbirdio/kubernetes-operator/internal/netbirdmock"
"github.com/netbirdio/netbird/shared/management/http/api"
)
var _ = Describe("NetworkRouter Controller", func() {
Context("When reconciling a resource", func() {
ctx := context.Background()
var netRouterRec *NetworkRouterReconciler
var setupKeyRec *SetupKeyReconciler
var groupRec *GroupReconciler
nn := client.ObjectKey{
Name: "test-resource",
Namespace: "network-router",
}
BeforeEach(func() {
nbClient := netbirdmock.Client()
netRouterRec = &NetworkRouterReconciler{
Client: k8sClient,
Netbird: nbClient,
ClientImage: "docker.io/netbirdio/netbird:latest",
ManagementURL: "https://netbird.io",
}
setupKeyRec = &SetupKeyReconciler{
Client: k8sClient,
Netbird: nbClient,
}
groupRec = &GroupReconciler{
Client: k8sClient,
Netbird: nbClient,
}
ns := &corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Namespace,
},
}
Expect(k8sClient.Create(ctx, ns)).To(Succeed())
})
AfterEach(func() {
ns := &corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Namespace,
},
}
err := k8sClient.Get(ctx, client.ObjectKeyFromObject(ns), ns)
if kerrors.IsNotFound(err) {
return
}
Expect(err).ToNot(HaveOccurred())
Expect(k8sClient.Delete(ctx, ns)).To(Succeed())
})
It("creates a routing peer along with a deployment", func() {
zoneReq := api.ZoneRequest{
Name: "cluster.local",
Domain: "cluster.local",
}
_, err := netRouterRec.Netbird.DNSZones.CreateZone(ctx, zoneReq)
Expect(err).ToNot(HaveOccurred())
netRouter := &nbv1alpha1.NetworkRouter{
ObjectMeta: metav1.ObjectMeta{
Name: nn.Name,
Namespace: nn.Namespace,
},
Spec: nbv1alpha1.NetworkRouterSpec{
DNSZoneRef: nbv1alpha1.DNSZoneReference{
Name: "cluster.local",
},
},
}
Expect(k8sClient.Create(ctx, netRouter)).To(Succeed())
group := &nbv1alpha1.Group{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("networkrouter-%s", netRouter.Name),
Namespace: nn.Namespace,
},
}
setupKey := &nbv1alpha1.SetupKey{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("networkrouter-%s", netRouter.Name),
Namespace: nn.Namespace,
},
}
dep := &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("networkrouter-%s", netRouter.Name),
Namespace: nn.Namespace,
},
}
for range 3 {
_, err := netRouterRec.Reconcile(ctx, reconcile.Request{NamespacedName: nn})
Expect(err).NotTo(HaveOccurred())
_, err = groupRec.Reconcile(ctx, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(group)})
Expect(err).NotTo(HaveOccurred())
_, err = setupKeyRec.Reconcile(ctx, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(setupKey)})
Expect(err).NotTo(HaveOccurred())
}
err = k8sClient.Get(ctx, nn, netRouter)
Expect(err).NotTo(HaveOccurred())
Expect(netRouter.Status.NetworkID).ToNot(BeEmpty())
_, err = netRouterRec.Netbird.Networks.Routers(netRouter.Status.NetworkID).Get(ctx, netRouter.Status.RoutingPeerID)
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(group), group)
Expect(err).NotTo(HaveOccurred())
Expect(group.OwnerReferences[0].UID).To(Equal(netRouter.UID))
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(setupKey), setupKey)
Expect(err).NotTo(HaveOccurred())
Expect(setupKey.OwnerReferences[0].UID).To(Equal(netRouter.UID))
err = k8sClient.Get(ctx, client.ObjectKeyFromObject(dep), dep)
Expect(err).NotTo(HaveOccurred())
Expect(dep.OwnerReferences[0].UID).To(Equal(netRouter.UID))
routingPeerResp, err := netRouterRec.Netbird.Networks.Routers(netRouter.Status.NetworkID).Get(ctx, netRouter.Status.RoutingPeerID)
Expect(err).NotTo(HaveOccurred())
Expect((*routingPeerResp.PeerGroups)[0]).To(Equal(group.Status.GroupID))
})
})
})
+5 -29
View File
@@ -3,7 +3,6 @@ package controller
import (
"context"
"errors"
"fmt"
"time"
"github.com/fluxcd/pkg/runtime/conditions"
@@ -12,7 +11,6 @@ import (
"github.com/netbirdio/netbird/shared/management/http/api"
corev1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
corev1ac "k8s.io/client-go/applyconfigurations/core/v1"
"k8s.io/utils/ptr"
ctrl "sigs.k8s.io/controller-runtime"
@@ -20,6 +18,7 @@ import (
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
"github.com/netbirdio/kubernetes-operator/internal/netbirdutil"
"github.com/netbirdio/kubernetes-operator/internal/ssautil"
)
@@ -52,32 +51,9 @@ func (r *SetupKeyReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c
return r.reconcileDelete(ctx, sp, setupKey)
}
// Get ids for auto groups.
autoGroupIDs := []string{}
for _, ref := range setupKey.Spec.AutoGroups {
switch {
case ref.ID != nil:
_, err := r.Netbird.Groups.Get(ctx, *ref.ID)
if err != nil {
return ctrl.Result{}, err
}
autoGroupIDs = append(autoGroupIDs, *ref.ID)
case ref.LocalRef != nil:
group := nbv1alpha1.Group{
ObjectMeta: metav1.ObjectMeta{
Name: ref.LocalRef.Name,
Namespace: setupKey.Namespace,
},
}
err = r.Client.Get(ctx, client.ObjectKeyFromObject(&group), &group)
if err != nil {
return ctrl.Result{}, err
}
if group.Status.GroupID == "" {
return ctrl.Result{}, fmt.Errorf("group %s in auto groups list is not ready", group.Name)
}
autoGroupIDs = append(autoGroupIDs, group.Status.GroupID)
}
autoGroupIDs, err := netbirdutil.GetGroupIDs(ctx, r.Client, r.Netbird, setupKey.Spec.AutoGroups, setupKey.Namespace)
if err != nil {
return ctrl.Result{}, err
}
controllerutil.AddFinalizer(setupKey, nbv1alpha1.NetbirdFinalizer)
@@ -151,7 +127,7 @@ func (r *SetupKeyReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c
AutoGroups: autoGroupIDs,
Ephemeral: ptr.To(setupKey.Spec.Ephemeral),
ExpiresIn: expiresIn,
Name: req.Name,
Name: setupKey.Spec.Name,
Type: "reusable",
UsageLimit: 0,
}
@@ -50,6 +50,9 @@ var _ = Describe("SetupKey Controller", func() {
Name: nn.Name,
Namespace: nn.Namespace,
},
Spec: nbv1alpha1.SetupKeySpec{
Name: "test",
},
}
Expect(k8sClient.Create(ctx, setupKey)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{NamespacedName: nn})
@@ -80,6 +83,9 @@ var _ = Describe("SetupKey Controller", func() {
Name: nn.Name,
Namespace: nn.Namespace,
},
Spec: nbv1alpha1.SetupKeySpec{
Name: "test",
},
}
Expect(k8sClient.Create(ctx, setupKey)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{NamespacedName: nn})
+66 -7
View File
@@ -31,6 +31,46 @@ func Client() *netbird.Client {
}
return output
})
addHandler(mux, "networks", func(id string, input api.NetworkRequest, output api.Network) api.Network {
output.Id = id
output.Name = input.Name
output.Description = input.Description
return output
})
addHandler(mux, "networks/{network}/routers", func(id string, input api.NetworkRouterRequest, output api.NetworkRouter) api.NetworkRouter {
output.Id = id
output.Enabled = input.Enabled
output.Masquerade = input.Masquerade
output.Metric = input.Metric
output.PeerGroups = input.PeerGroups
return output
})
addHandler(mux, "networks/{network}/resources", func(id string, input api.NetworkResourceRequest, output api.NetworkResource) api.NetworkResource {
output.Id = id
output.Address = input.Address
output.Description = input.Description
output.Enabled = input.Enabled
return output
})
addHandler(mux, "dns/zones", func(id string, input api.ZoneRequest, output api.Zone) api.Zone {
output.Id = id
output.Name = input.Name
output.Domain = input.Domain
output.DistributionGroups = input.DistributionGroups
output.EnableSearchDomain = input.EnableSearchDomain
if input.Enabled != nil {
output.Enabled = *input.Enabled
}
return output
})
addHandler(mux, "dns/zones/{zone}/records", func(id string, input api.DNSRecordRequest, output api.DNSRecord) api.DNSRecord {
output.Id = id
output.Name = input.Name
output.Ttl = input.Ttl
output.Type = input.Type
output.Content = input.Content
return output
})
srv := httptest.NewServer(mux)
return netbird.New(srv.URL, "ABC")
@@ -38,14 +78,33 @@ func Client() *netbird.Client {
func addHandler[T, U any](mux *http.ServeMux, resource string, convertFn func(string, U, T) T) {
var itemMx sync.RWMutex
items := map[string]T{}
store := map[string]T{}
mux.Handle(fmt.Sprintf("GET /api/%s", resource), http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) {
itemMx.RLock()
defer itemMx.RUnlock()
items := make([]T, 0, len(store))
for _, v := range store {
items = append(items, v)
}
b, err := json.Marshal(items)
if err != nil {
util.WriteErrorResponse("Marshal Error", http.StatusInternalServerError, rw)
return
}
_, err = rw.Write(b)
if err != nil {
util.WriteErrorResponse("Write Error", http.StatusInternalServerError, rw)
return
}
}))
mux.Handle(fmt.Sprintf("GET /api/%s/{id}", resource), http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) {
itemMx.RLock()
defer itemMx.RUnlock()
id := req.PathValue("id")
respData, ok := items[id]
respData, ok := store[id]
if !ok {
util.WriteErrorResponse("Not Found", http.StatusNotFound, rw)
return
@@ -79,7 +138,7 @@ func addHandler[T, U any](mux *http.ServeMux, resource string, convertFn func(st
id := fmt.Sprintf("id-%d", rand.Int64())
var zero T
respData := convertFn(id, reqData, zero)
items[id] = respData
store[id] = respData
b, err = json.Marshal(respData)
if err != nil {
util.WriteErrorResponse("Marshal Error", http.StatusInternalServerError, rw)
@@ -96,7 +155,7 @@ func addHandler[T, U any](mux *http.ServeMux, resource string, convertFn func(st
defer itemMx.Unlock()
id := req.PathValue("id")
respData, ok := items[id]
respData, ok := store[id]
if !ok {
util.WriteErrorResponse("Not Found", http.StatusNotFound, rw)
return
@@ -114,7 +173,7 @@ func addHandler[T, U any](mux *http.ServeMux, resource string, convertFn func(st
return
}
respData = convertFn(id, reqData, respData)
items[id] = respData
store[id] = respData
b, err = json.Marshal(respData)
if err != nil {
util.WriteErrorResponse("Marshal Error", http.StatusInternalServerError, rw)
@@ -131,11 +190,11 @@ func addHandler[T, U any](mux *http.ServeMux, resource string, convertFn func(st
defer itemMx.Unlock()
id := req.PathValue("id")
_, ok := items[id]
_, ok := store[id]
if !ok {
util.WriteErrorResponse("Not Found", http.StatusNotFound, rw)
return
}
delete(items, id)
delete(store, id)
}))
}
+42
View File
@@ -0,0 +1,42 @@
package netbirdutil
import (
"context"
"fmt"
netbird "github.com/netbirdio/netbird/shared/management/client/rest"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
nbv1alpha1 "github.com/netbirdio/kubernetes-operator/api/v1alpha1"
)
func GetGroupIDs(ctx context.Context, k8sClient client.Client, nbClient *netbird.Client, refs []nbv1alpha1.ResourceReference, namespace string) ([]string, error) {
groupIDs := []string{}
for _, ref := range refs {
switch {
case ref.ID != nil:
_, err := nbClient.Groups.Get(ctx, *ref.ID)
if err != nil {
return nil, err
}
groupIDs = append(groupIDs, *ref.ID)
case ref.LocalRef != nil:
group := nbv1alpha1.Group{
ObjectMeta: metav1.ObjectMeta{
Name: ref.LocalRef.Name,
Namespace: namespace,
},
}
err := k8sClient.Get(ctx, client.ObjectKeyFromObject(&group), &group)
if err != nil {
return nil, err
}
if group.Status.GroupID == "" {
return nil, fmt.Errorf("group %s in groups list is not ready", group.Name)
}
groupIDs = append(groupIDs, group.Status.GroupID)
}
}
return groupIDs, nil
}
+24
View File
@@ -0,0 +1,24 @@
package netbirdutil
import (
"context"
"fmt"
"slices"
netbird "github.com/netbirdio/netbird/shared/management/client/rest"
"github.com/netbirdio/netbird/shared/management/http/api"
)
func GetDNSZoneByName(ctx context.Context, nbClient *netbird.Client, name string) (api.Zone, error) {
resp, err := nbClient.DNSZones.ListZones(ctx)
if err != nil {
return api.Zone{}, err
}
zoneIdx := slices.IndexFunc(resp, func(zone api.Zone) bool {
return zone.Name == name
})
if zoneIdx == -1 {
return api.Zone{}, fmt.Errorf("zone with name %s cannot be found", name)
}
return resp[zoneIdx], nil
}