Add unit tests to new controllers and fix minor bugs (#12)

This commit is contained in:
M. Essam
2025-03-28 08:55:41 +01:00
committed by GitHub
parent ac8348cda6
commit 6a33bffb65
22 changed files with 3788 additions and 313 deletions
+327 -30
View File
@@ -2,16 +2,21 @@ package controller
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"k8s.io/apimachinery/pkg/api/errors"
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
netbirdiov1 "github.com/netbirdio/kubernetes-operator/api/v1"
"github.com/netbirdio/kubernetes-operator/internal/util"
netbird "github.com/netbirdio/netbird/management/client/rest"
"github.com/netbirdio/netbird/management/server/http/api"
)
var _ = Describe("NBGroup Controller", func() {
@@ -22,49 +27,341 @@ var _ = Describe("NBGroup Controller", func() {
typeNamespacedName := types.NamespacedName{
Name: resourceName,
Namespace: "default", // TODO(user):Modify as needed
Namespace: "default",
}
nbgroup := &netbirdiov1.NBGroup{}
var netbirdClient *netbird.Client
var mux *http.ServeMux
var server *httptest.Server
var nbGroup netbirdiov1.NBGroup
BeforeEach(func() {
Skip("Not implemented yet")
By("creating the custom resource for the Kind NBGroup")
err := k8sClient.Get(ctx, typeNamespacedName, nbgroup)
if err != nil && errors.IsNotFound(err) {
resource := &netbirdiov1.NBGroup{
ObjectMeta: metav1.ObjectMeta{
Name: resourceName,
Namespace: "default",
mux = &http.ServeMux{}
server = httptest.NewServer(mux)
netbirdClient = netbird.New(server.URL, "ABC")
err := k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
if err == nil {
deleteErr := k8sClient.Delete(ctx, &nbGroup)
Expect(deleteErr).NotTo(HaveOccurred())
}
if err == nil || errors.IsNotFound(err) {
nbGroup = netbirdiov1.NBGroup{
ObjectMeta: v1.ObjectMeta{
Name: resourceName,
Namespace: typeNamespacedName.Namespace,
Finalizers: []string{"netbird.io/group-cleanup"},
},
Spec: netbirdiov1.NBGroupSpec{
Name: resourceName,
},
// TODO(user): Specify other spec details if needed.
}
Expect(k8sClient.Create(ctx, resource)).To(Succeed())
err = k8sClient.Create(ctx, &nbGroup)
Expect(err).NotTo(HaveOccurred())
}
})
AfterEach(func() {
// TODO(user): Cleanup logic after each test, like removing the resource instance.
server.Close()
resource := &netbirdiov1.NBGroup{}
err := k8sClient.Get(ctx, typeNamespacedName, resource)
if errors.IsNotFound(err) {
return
}
Expect(err).NotTo(HaveOccurred())
By("Cleanup the specific resource instance NBGroup")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
})
It("should successfully reconcile the resource", func() {
Skip("Not implemented yet")
By("Reconciling the created resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
if len(resource.Finalizers) > 0 {
resource.Finalizers = nil
Expect(k8sClient.Update(ctx, resource)).To(Succeed())
}
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
if resource.DeletionTimestamp == nil {
By("Cleanup the specific resource instance NBGroup")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
}
})
When("Group doesn't exist", func() {
It("should create group", func() {
By("Reconciling the created resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodGet {
_, err := w.Write([]byte("[]"))
Expect(err).NotTo(HaveOccurred())
} else {
resp := api.Group{
Id: "Test",
Name: resourceName,
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(err).NotTo(HaveOccurred())
Expect(nbGroup.Status.GroupID).NotTo(BeNil())
Expect(*nbGroup.Status.GroupID).To(Equal("Test"))
Expect(nbGroup.Status.Conditions).To(HaveLen(1))
Expect(nbGroup.Status.Conditions[0].Status).To(BeEquivalentTo(v1.ConditionTrue))
Expect(nbGroup.Status.Conditions[0].Type).To(Equal(netbirdiov1.NBSetupKeyReady))
})
})
When("Group already exists", func() {
It("should use existing group", func() {
By("Reconciling the created resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "Test",
Name: resourceName,
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(err).NotTo(HaveOccurred())
Expect(nbGroup.Status.GroupID).NotTo(BeNil())
Expect(*nbGroup.Status.GroupID).To(Equal("Test"))
Expect(nbGroup.Status.Conditions).To(HaveLen(1))
Expect(nbGroup.Status.Conditions[0].Status).To(BeEquivalentTo(v1.ConditionTrue))
Expect(nbGroup.Status.Conditions[0].Type).To(Equal(netbirdiov1.NBSetupKeyReady))
})
})
When("NBGroup is set for deletion", func() {
deleteGroup := func() {
GinkgoHelper()
By("Adding the group ID in status")
nbGroup.Status.GroupID = util.Ptr("Test")
err := k8sClient.Status().Update(ctx, &nbGroup)
Expect(err).NotTo(HaveOccurred())
By("Deleting the object")
err = k8sClient.Delete(ctx, &nbGroup)
Expect(err).NotTo(HaveOccurred())
}
When("Group is not linked to any resources", func() {
It("should delete group", func() {
deleteGroup()
By("Reconciling the deleting resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
method := ""
mux.HandleFunc("/api/groups/Test", func(w http.ResponseWriter, r *http.Request) {
method = r.Method
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(errors.IsNotFound(err)).To(BeTrue())
Expect(method).To(Equal(http.MethodDelete))
})
})
When("Group is linked to other resources", func() {
It("should return error", func() {
deleteGroup()
By("Reconciling the deleting resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
method := ""
mux.HandleFunc("/api/groups/Test", func(w http.ResponseWriter, r *http.Request) {
method = r.Method
w.WriteHeader(400)
_, err := w.Write([]byte(`{"message": "group has been linked to Policy: meow"}`))
Expect(err).NotTo(HaveOccurred())
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).To(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(errors.IsNotFound(err)).To(BeFalse())
Expect(method).To(Equal(http.MethodDelete))
})
})
When("Group already exists in another namespace", func() {
It("Should delete NBGroup after linked failure", func() {
deleteGroup()
otherGroup := &netbirdiov1.NBGroup{
ObjectMeta: v1.ObjectMeta{
Name: nbGroup.Name,
Namespace: "kube-system",
},
Spec: netbirdiov1.NBGroupSpec{
Name: nbGroup.Spec.Name,
},
}
Expect(k8sClient.Create(ctx, otherGroup)).To(Succeed())
otherGroup.Status.GroupID = nbGroup.Status.GroupID
Expect(k8sClient.Status().Update(ctx, otherGroup)).To(Succeed())
By("Reconciling the deleting resource")
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
method := ""
mux.HandleFunc("/api/groups/Test", func(w http.ResponseWriter, r *http.Request) {
method = r.Method
w.WriteHeader(400)
_, err := w.Write([]byte(`{"message": "group has been linked to Policy: meow"}`))
Expect(err).NotTo(HaveOccurred())
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(errors.IsNotFound(err)).To(BeTrue())
Expect(method).To(Equal(http.MethodDelete))
})
})
})
When("Group already exists with different ID", func() {
It("should re-use existing group ID", func() {
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "Test",
Name: resourceName,
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
nbGroup.Status.GroupID = util.Ptr("Toast")
Expect(k8sClient.Status().Update(ctx, &nbGroup)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(err).NotTo(HaveOccurred())
Expect(nbGroup.Status.GroupID).NotTo(BeNil())
Expect(*nbGroup.Status.GroupID).To(Equal("Test"))
Expect(nbGroup.Status.Conditions).To(HaveLen(1))
Expect(nbGroup.Status.Conditions[0].Status).To(BeEquivalentTo(v1.ConditionTrue))
Expect(nbGroup.Status.Conditions[0].Type).To(Equal(netbirdiov1.NBSetupKeyReady))
})
})
When("Group deleted from NetBird API", func() {
It("Should requeue and create group on next run", func() {
controllerReconciler := &NBGroupReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodGet {
resp := []api.Group{}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
if r.Method == http.MethodPost {
resp := api.Group{
Id: "Test",
Name: resourceName,
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
nbGroup.Status.GroupID = util.Ptr("Toast")
Expect(k8sClient.Status().Update(ctx, &nbGroup)).To(Succeed())
res, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(res.Requeue).To(BeTrue())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(err).NotTo(HaveOccurred())
Expect(nbGroup.Status.GroupID).To(BeNil())
Expect(nbGroup.Status.Conditions).To(HaveLen(1))
Expect(nbGroup.Status.Conditions[0].Status).To(BeEquivalentTo(v1.ConditionFalse))
_, err = controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
err = k8sClient.Get(ctx, typeNamespacedName, &nbGroup)
Expect(err).NotTo(HaveOccurred())
Expect(nbGroup.Status.GroupID).NotTo(BeNil())
Expect(*nbGroup.Status.GroupID).To(Equal("Test"))
Expect(nbGroup.Status.Conditions).To(HaveLen(1))
Expect(nbGroup.Status.Conditions[0].Status).To(BeEquivalentTo(v1.ConditionTrue))
})
Expect(err).NotTo(HaveOccurred())
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
// Example: If you expect a certain status condition after reconciliation, verify it here.
})
})
})
+23 -15
View File
@@ -35,6 +35,11 @@ var (
errNetBirdAPI = fmt.Errorf("netbird API error")
)
const (
protocolTCP = "tcp"
protocolUDP = "udp"
)
// getResources get all NBResource objects in policy.status.managedServiceList
func (r *NBPolicyReconciler) getResources(ctx context.Context, nbPolicy *netbirdiov1.NBPolicy, logger logr.Logger) ([]netbirdiov1.NBResource, error) {
var resourceList []netbirdiov1.NBResource
@@ -63,8 +68,8 @@ func (r *NBPolicyReconciler) getResources(ctx context.Context, nbPolicy *netbird
// returns map[protocol] => ports, destination group IDs
func (r *NBPolicyReconciler) mapResources(ctx context.Context, nbPolicy *netbirdiov1.NBPolicy, resources []netbirdiov1.NBResource, logger logr.Logger) (map[string][]int32, []string, error) {
portMapping := map[string]map[int32]interface{}{
"tcp": make(map[int32]interface{}),
"udp": make(map[int32]interface{}),
protocolTCP: make(map[int32]interface{}),
protocolUDP: make(map[int32]interface{}),
}
groups, err := r.groupNamesToIDs(ctx, nbPolicy.Spec.DestinationGroups, logger)
if err != nil {
@@ -77,16 +82,17 @@ func (r *NBPolicyReconciler) mapResources(ctx context.Context, nbPolicy *netbird
groups = append(groups, resource.Status.Groups...)
for _, p := range resource.Spec.TCPPorts {
portMapping["tcp"][p] = nil
portMapping[protocolTCP][p] = nil
}
for _, p := range resource.Spec.UDPPorts {
portMapping["udp"][p] = nil
portMapping[protocolUDP][p] = nil
}
}
}
ports := make(map[string][]int32)
for k, vs := range portMapping {
ports[k] = nil
for v := range vs {
ports[k] = append(ports[k], v)
}
@@ -131,7 +137,7 @@ func (r *NBPolicyReconciler) createPolicy(ctx context.Context, nbPolicy *netbird
func (r *NBPolicyReconciler) updatePolicy(ctx context.Context, policyID *string, nbPolicy *netbirdiov1.NBPolicy, protocol string, sourceGroupIDs, destinationGroupIDs, ports []string, logger logr.Logger) (*string, bool, error) {
policyName := fmt.Sprintf("%s %s", nbPolicy.Spec.Name, strings.ToUpper(protocol))
logger.Info("Updating NetBird Policy", "name", policyName, "description", nbPolicy.Spec.Description, "protocol", protocol, "sources", sourceGroupIDs, "destinations", destinationGroupIDs, "ports", ports, "bidirectional", nbPolicy.Spec.Bidirectional)
policy, err := r.netbird.Policies.Update(ctx, *policyID, api.PutApiPoliciesPolicyIdJSONRequestBody{
_, err := r.netbird.Policies.Update(ctx, *policyID, api.PutApiPoliciesPolicyIdJSONRequestBody{
Enabled: true,
Name: policyName,
Description: &nbPolicy.Spec.Description,
@@ -163,11 +169,10 @@ func (r *NBPolicyReconciler) updatePolicy(ctx context.Context, policyID *string,
policyID = nil
requeue = true
nbPolicy.Status.Conditions = netbirdiov1.NBConditionFalse("Gone", "Policy deleted from NetBird API")
} else if err != nil {
return nil, false, err
}
if err == nil && (policyID == nil || *policy.Id != *policyID) {
policyID = policy.Id
}
return policyID, requeue, nil
}
@@ -192,6 +197,9 @@ func (r *NBPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Request) (r
originalPolicy := nbPolicy.DeepCopy()
defer func() {
if originalPolicy.DeletionTimestamp != nil && len(nbPolicy.Finalizers) == 0 {
return
}
if !originalPolicy.Status.Equal(nbPolicy.Status) {
updateErr := r.Client.Status().Update(ctx, &nbPolicy)
if updateErr != nil {
@@ -207,7 +215,7 @@ func (r *NBPolicyReconciler) Reconcile(ctx context.Context, req ctrl.Request) (r
if len(nbPolicy.Finalizers) == 0 {
return ctrl.Result{}, nil
}
return ctrl.Result{}, r.handleDelete(ctx, nbPolicy, logger)
return ctrl.Result{}, r.handleDelete(ctx, &nbPolicy, logger)
}
resourceList, err := r.getResources(ctx, &nbPolicy, logger)
@@ -244,9 +252,9 @@ func (r *NBPolicyReconciler) syncPolicy(ctx context.Context, nbPolicy *netbirdio
for protocol, ports := range portMapping {
var policyID *string
switch protocol {
case "tcp":
case protocolTCP:
policyID = nbPolicy.Status.TCPPolicyID
case "udp":
case protocolUDP:
policyID = nbPolicy.Status.UDPPolicyID
default:
logger.Error(errKubernetesAPI, "Unknown protocol", "protocol", protocol)
@@ -306,9 +314,9 @@ func (r *NBPolicyReconciler) syncPolicy(ctx context.Context, nbPolicy *netbirdio
}
switch protocol {
case "tcp":
case protocolTCP:
nbPolicy.Status.TCPPolicyID = policyID
case "udp":
case protocolUDP:
nbPolicy.Status.UDPPolicyID = policyID
default:
logger.Error(errKubernetesAPI, "Unknown protocol", "protocol", protocol)
@@ -320,7 +328,7 @@ func (r *NBPolicyReconciler) syncPolicy(ctx context.Context, nbPolicy *netbirdio
return requeue, nil
}
func (r *NBPolicyReconciler) handleDelete(ctx context.Context, nbPolicy netbirdiov1.NBPolicy, logger logr.Logger) error {
func (r *NBPolicyReconciler) handleDelete(ctx context.Context, nbPolicy *netbirdiov1.NBPolicy, logger logr.Logger) error {
if nbPolicy.Status.TCPPolicyID != nil {
err := r.netbird.Policies.Delete(ctx, *nbPolicy.Status.TCPPolicyID)
if err != nil && !strings.Contains("not found", err.Error()) {
@@ -337,7 +345,7 @@ func (r *NBPolicyReconciler) handleDelete(ctx context.Context, nbPolicy netbirdi
}
if util.Contains(nbPolicy.Finalizers, "netbird.io/cleanup") {
nbPolicy.Finalizers = util.Without(nbPolicy.Finalizers, "netbird.io/cleanup")
err := r.Client.Update(ctx, &nbPolicy)
err := r.Client.Update(ctx, nbPolicy)
if err != nil {
logger.Error(errKubernetesAPI, "Error updating NBPolicy", "err", err)
return err
+600 -23
View File
@@ -2,7 +2,12 @@ package controller
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"github.com/go-logr/logr"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"k8s.io/apimachinery/pkg/api/errors"
@@ -12,59 +17,631 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
netbirdiov1 "github.com/netbirdio/kubernetes-operator/api/v1"
"github.com/netbirdio/kubernetes-operator/internal/util"
netbird "github.com/netbirdio/netbird/management/client/rest"
"github.com/netbirdio/netbird/management/server/http/api"
ctrl "sigs.k8s.io/controller-runtime"
)
var _ = Describe("NBPolicy Controller", func() {
Context("When reconciling a resource", func() {
const resourceName = "test-resource"
var resourceName = "test-resource"
ctx := context.Background()
typeNamespacedName := types.NamespacedName{
Name: resourceName,
Namespace: "default", // TODO(user):Modify as needed
Name: resourceName,
}
nbpolicy := &netbirdiov1.NBPolicy{}
var netbirdClient *netbird.Client
var mux *http.ServeMux
var server *httptest.Server
BeforeEach(func() {
Skip("Not implemented yet")
ctrl.SetLogger(logr.New(GinkgoLogr.GetSink()))
mux = &http.ServeMux{}
server = httptest.NewServer(mux)
netbirdClient = netbird.New(server.URL, "ABC")
By("creating the custom resource for the Kind NBPolicy")
err := k8sClient.Get(ctx, typeNamespacedName, nbpolicy)
if err != nil && errors.IsNotFound(err) {
resource := &netbirdiov1.NBPolicy{
ObjectMeta: metav1.ObjectMeta{
Name: resourceName,
Namespace: "default",
Name: resourceName,
Finalizers: []string{"netbird.io/cleanup"},
},
Spec: netbirdiov1.NBPolicySpec{
Name: "Test",
SourceGroups: []string{"All"},
Bidirectional: true,
},
// TODO(user): Specify other spec details if needed.
}
Expect(k8sClient.Create(ctx, resource)).To(Succeed())
nbpolicy = resource
}
})
AfterEach(func() {
// TODO(user): Cleanup logic after each test, like removing the resource instance.
resource := &netbirdiov1.NBPolicy{}
err := k8sClient.Get(ctx, typeNamespacedName, resource)
Expect(err).NotTo(HaveOccurred())
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
By("Cleanup the specific resource instance NBPolicy")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
})
It("should successfully reconcile the resource", func() {
Skip("Not implemented yet")
By("Reconciling the created resource")
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
if len(resource.Finalizers) > 0 {
resource.Finalizers = nil
Expect(k8sClient.Update(ctx, resource)).To(Succeed())
}
By("Cleanup the specific resource instance NBPolicy")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
}
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
nbresource := &netbirdiov1.NBResource{}
err = k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "test"}, nbresource)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
By("Cleanup the specific resource instance NBResource")
Expect(k8sClient.Delete(ctx, nbresource)).To(Succeed())
}
})
When("Not enough information to create policy", func() {
It("should not create any policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
})
})
When("Enough information to create TCP policy", func() {
It("should create 1 policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbResource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "meow",
Groups: []string{"test"},
NetworkID: "test",
Address: "test.default.svc.cluster.local",
PolicyName: resourceName,
TCPPorts: []int32{443},
},
}
Expect(k8sClient.Create(ctx, nbResource)).To(Succeed())
nbResource.Status = netbirdiov1.NBResourceStatus{
TCPPorts: []int32{443},
PolicyName: &resourceName,
Groups: []string{"test"},
}
Expect(k8sClient.Status().Update(ctx, nbResource)).To(Succeed())
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/test")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyCreated := false
mux.HandleFunc("/api/policies", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodPost {
var policyReq api.PostApiPoliciesJSONRequestBody
bs, err := io.ReadAll(r.Body)
Expect(err).NotTo(HaveOccurred())
err = json.Unmarshal(bs, &policyReq)
Expect(err).NotTo(HaveOccurred())
Expect(policyReq.Name).To(Equal("Test TCP"))
Expect(policyReq.Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Enabled).To(BeTrue())
Expect(policyReq.SourcePostureChecks).To(BeNil())
Expect(policyReq.Rules).To(HaveLen(1))
Expect(policyReq.Rules[0].Action).To(BeEquivalentTo(api.PolicyRuleActionAccept))
Expect(policyReq.Rules[0].Bidirectional).To(BeTrue())
Expect(policyReq.Rules[0].Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Rules[0].DestinationResource).To(BeNil())
Expect(policyReq.Rules[0].Destinations).NotTo(BeNil())
Expect(*policyReq.Rules[0].Destinations).To(HaveLen(1))
Expect((*policyReq.Rules[0].Destinations)[0]).To(Equal("test"))
Expect(policyReq.Rules[0].Enabled).To(BeTrue())
Expect(policyReq.Rules[0].Name).To(Equal("Test TCP"))
Expect(policyReq.Rules[0].Ports).NotTo(BeNil())
Expect((*policyReq.Rules[0].Ports)).To(HaveLen(1))
Expect((*policyReq.Rules[0].Ports)[0]).To(Equal("443"))
Expect(policyReq.Rules[0].Protocol).To(BeEquivalentTo(api.PolicyRuleProtocolTcp))
Expect(policyReq.Rules[0].SourceResource).To(BeNil())
Expect(policyReq.Rules[0].Sources).NotTo(BeNil())
Expect(*policyReq.Rules[0].Sources).To(HaveLen(1))
Expect((*policyReq.Rules[0].Sources)[0]).To(Equal("meow"))
policyCreated = true
resp := api.Policy{
Id: &resourceName,
}
bs, err = json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyCreated).To(BeTrue())
})
})
When("TCP information no longer sufficient", func() {
It("should delete tcp policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/noexist")
nbpolicy.Status.TCPPolicyID = util.Ptr("policyid")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyDeleted := false
mux.HandleFunc("/api/policies/policyid", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodDelete {
policyDeleted = true
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyDeleted).To(BeTrue())
})
})
When("Enough information to create UDP policy", func() {
It("should create 1 policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbResource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "meow",
Groups: []string{"test"},
NetworkID: "test",
Address: "test.default.svc.cluster.local",
PolicyName: resourceName,
UDPPorts: []int32{443},
},
}
Expect(k8sClient.Create(ctx, nbResource)).To(Succeed())
nbResource.Status = netbirdiov1.NBResourceStatus{
UDPPorts: []int32{443},
PolicyName: &resourceName,
Groups: []string{"test"},
}
Expect(k8sClient.Status().Update(ctx, nbResource)).To(Succeed())
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/test")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyCreated := false
mux.HandleFunc("/api/policies", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodPost {
var policyReq api.PostApiPoliciesJSONRequestBody
bs, err := io.ReadAll(r.Body)
Expect(err).NotTo(HaveOccurred())
err = json.Unmarshal(bs, &policyReq)
Expect(err).NotTo(HaveOccurred())
Expect(policyReq.Name).To(Equal("Test UDP"))
Expect(policyReq.Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Enabled).To(BeTrue())
Expect(policyReq.SourcePostureChecks).To(BeNil())
Expect(policyReq.Rules).To(HaveLen(1))
Expect(policyReq.Rules[0].Action).To(BeEquivalentTo(api.PolicyRuleActionAccept))
Expect(policyReq.Rules[0].Bidirectional).To(BeTrue())
Expect(policyReq.Rules[0].Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Rules[0].DestinationResource).To(BeNil())
Expect(policyReq.Rules[0].Destinations).NotTo(BeNil())
Expect(*policyReq.Rules[0].Destinations).To(HaveLen(1))
Expect((*policyReq.Rules[0].Destinations)[0]).To(Equal("test"))
Expect(policyReq.Rules[0].Enabled).To(BeTrue())
Expect(policyReq.Rules[0].Name).To(Equal("Test UDP"))
Expect(policyReq.Rules[0].Ports).NotTo(BeNil())
Expect((*policyReq.Rules[0].Ports)).To(HaveLen(1))
Expect((*policyReq.Rules[0].Ports)[0]).To(Equal("443"))
Expect(policyReq.Rules[0].Protocol).To(BeEquivalentTo(api.PolicyRuleProtocolUdp))
Expect(policyReq.Rules[0].SourceResource).To(BeNil())
Expect(policyReq.Rules[0].Sources).NotTo(BeNil())
Expect(*policyReq.Rules[0].Sources).To(HaveLen(1))
Expect((*policyReq.Rules[0].Sources)[0]).To(Equal("meow"))
policyCreated = true
resp := api.Policy{
Id: &resourceName,
}
bs, err = json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyCreated).To(BeTrue())
})
})
When("UDP information no longer sufficient", func() {
It("should delete udp policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/noexist")
nbpolicy.Status.UDPPolicyID = util.Ptr("policyid")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyDeleted := false
mux.HandleFunc("/api/policies/policyid", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodDelete {
policyDeleted = true
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyDeleted).To(BeTrue())
})
})
When("Existing protocol gets restricted", func() {
It("Should delete protocol policy", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbResource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "meow",
Groups: []string{"test"},
NetworkID: "test",
Address: "test.default.svc.cluster.local",
PolicyName: resourceName,
TCPPorts: []int32{443},
},
}
Expect(k8sClient.Create(ctx, nbResource)).To(Succeed())
nbResource.Status = netbirdiov1.NBResourceStatus{
TCPPorts: []int32{443},
PolicyName: &resourceName,
Groups: []string{"test"},
}
Expect(k8sClient.Status().Update(ctx, nbResource)).To(Succeed())
nbpolicy.Spec.Protocols = []string{"udp"}
Expect(k8sClient.Update(ctx, nbpolicy)).To(Succeed())
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/test")
nbpolicy.Status.TCPPolicyID = util.Ptr("policyid")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyDeleted := false
mux.HandleFunc("/api/policies/policyid", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodDelete {
policyDeleted = true
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyDeleted).To(BeTrue())
})
})
When("Updating existing policy", func() {
AfterEach(func() {
nbresource := &netbirdiov1.NBResource{}
err := k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "test-b"}, nbresource)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
By("Cleanup the specific resource instance NBResource")
Expect(k8sClient.Delete(ctx, nbresource)).To(Succeed())
}
})
It("Should give all information to Update method", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbResource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "meow",
Groups: []string{"test"},
NetworkID: "test",
Address: "test.default.svc.cluster.local",
PolicyName: resourceName,
TCPPorts: []int32{443},
},
}
Expect(k8sClient.Create(ctx, nbResource)).To(Succeed())
nbResource.Status = netbirdiov1.NBResourceStatus{
TCPPorts: []int32{443},
PolicyName: &resourceName,
Groups: []string{"test"},
}
Expect(k8sClient.Status().Update(ctx, nbResource)).To(Succeed())
nbResourceB := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "test-b",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "meow-b",
Groups: []string{"test-b"},
NetworkID: "test",
Address: "test-b.default.svc.cluster.local",
PolicyName: resourceName,
TCPPorts: []int32{80},
},
}
Expect(k8sClient.Create(ctx, nbResourceB)).To(Succeed())
nbResourceB.Status = netbirdiov1.NBResourceStatus{
TCPPorts: []int32{80},
PolicyName: &resourceName,
Groups: []string{"test-b"},
}
Expect(k8sClient.Status().Update(ctx, nbResourceB)).To(Succeed())
nbpolicy.Status.ManagedServiceList = append(nbpolicy.Status.ManagedServiceList, "default/test", "default/test-b")
nbpolicy.Status.TCPPolicyID = util.Ptr("policyid")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
resp := []api.Group{
{
Id: "meow",
Name: "All",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
policyUpdated := false
mux.HandleFunc("/api/policies/policyid", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodPut {
policyUpdated = true
var policyReq api.PostApiPoliciesJSONRequestBody
bs, err := io.ReadAll(r.Body)
Expect(err).NotTo(HaveOccurred())
err = json.Unmarshal(bs, &policyReq)
Expect(err).NotTo(HaveOccurred())
Expect(policyReq.Name).To(Equal("Test TCP"))
Expect(policyReq.Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Enabled).To(BeTrue())
Expect(policyReq.SourcePostureChecks).To(BeNil())
Expect(policyReq.Rules).To(HaveLen(1))
Expect(policyReq.Rules[0].Action).To(BeEquivalentTo(api.PolicyRuleActionAccept))
Expect(policyReq.Rules[0].Bidirectional).To(BeTrue())
Expect(policyReq.Rules[0].Description).To(Or(BeNil(), BeEquivalentTo(util.Ptr(""))))
Expect(policyReq.Rules[0].DestinationResource).To(BeNil())
Expect(policyReq.Rules[0].Destinations).NotTo(BeNil())
Expect(*policyReq.Rules[0].Destinations).To(HaveLen(2))
Expect((*policyReq.Rules[0].Destinations)).To(ConsistOf([]string{"test", "test-b"}))
Expect(policyReq.Rules[0].Enabled).To(BeTrue())
Expect(policyReq.Rules[0].Name).To(Equal("Test TCP"))
Expect(policyReq.Rules[0].Ports).NotTo(BeNil())
Expect((*policyReq.Rules[0].Ports)).To(HaveLen(2))
Expect((*policyReq.Rules[0].Ports)).To(ConsistOf([]string{"443", "80"}))
Expect(policyReq.Rules[0].Protocol).To(BeEquivalentTo(api.PolicyRuleProtocolTcp))
Expect(policyReq.Rules[0].SourceResource).To(BeNil())
Expect(policyReq.Rules[0].Sources).NotTo(BeNil())
Expect(*policyReq.Rules[0].Sources).To(HaveLen(1))
Expect((*policyReq.Rules[0].Sources)[0]).To(Equal("meow"))
_, err = w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(policyUpdated).To(BeTrue())
})
})
When("NBPolicy is set for deletion", func() {
It("should delete Policies", func() {
controllerReconciler := &NBPolicyReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
ClusterName: "Kubernetes",
}
nbpolicy.Status.TCPPolicyID = util.Ptr("policyidtcp")
nbpolicy.Status.UDPPolicyID = util.Ptr("policyidudp")
Expect(k8sClient.Status().Update(ctx, nbpolicy)).To(Succeed())
Expect(k8sClient.Delete(ctx, nbpolicy)).To(Succeed())
tcpPolicyDeleted := false
mux.HandleFunc("/api/policies/policyidtcp", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodDelete {
tcpPolicyDeleted = true
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
udpPolicyDeleted := false
mux.HandleFunc("/api/policies/policyidudp", func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodDelete {
udpPolicyDeleted = true
_, err := w.Write([]byte("{}"))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(tcpPolicyDeleted).To(BeTrue())
Expect(udpPolicyDeleted).To(BeTrue())
err = k8sClient.Get(ctx, typeNamespacedName, nbpolicy)
Expect(errors.IsNotFound(err)).To(BeTrue())
})
Expect(err).NotTo(HaveOccurred())
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
// Example: If you expect a certain status condition after reconciliation, verify it here.
})
})
})
+83 -10
View File
@@ -3,6 +3,7 @@ package controller
import (
"context"
"fmt"
"slices"
"strings"
"time"
@@ -48,6 +49,9 @@ func (r *NBResourceReconciler) Reconcile(ctx context.Context, req ctrl.Request)
originalResource := nbResource.DeepCopy()
defer func() {
if originalResource.DeletionTimestamp != nil && len(nbResource.Finalizers) == 0 {
return
}
if !originalResource.Status.Equal(nbResource.Status) {
updateErr := r.Client.Status().Update(ctx, nbResource)
if updateErr != nil {
@@ -110,8 +114,8 @@ func (r *NBResourceReconciler) handlePolicy(ctx context.Context, req ctrl.Reques
var nbPolicy netbirdiov1.NBPolicy
if nbResource.Spec.PolicyName == "" && nbResource.Status.PolicyName != nil {
// Remove self reference from policy status
nbResource.Status.PolicyName = nil
err := r.Client.Get(ctx, types.NamespacedName{Name: *nbResource.Status.PolicyName}, &nbPolicy)
nbResource.Status.PolicyName = nil
if err != nil {
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", nbResource.Spec.PolicyName)
return err
@@ -123,7 +127,24 @@ func (r *NBResourceReconciler) handlePolicy(ctx context.Context, req ctrl.Reques
}
} else {
// Update policy settings if any difference is found
// TODO: Handle updated policy name by removing reference from old policy name in status.policyName
if nbResource.Status.PolicyName != nil {
err := r.Client.Get(ctx, types.NamespacedName{Name: *nbResource.Status.PolicyName}, &nbPolicy)
if !errors.IsNotFound(err) {
if err != nil {
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", nbResource.Spec.PolicyName)
return err
}
if util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
nbPolicy.Status.ManagedServiceList = util.Without(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
err := r.Client.Status().Update(ctx, &nbPolicy)
if err != nil {
logger.Error(errKubernetesAPI, "error updating NBPolicy", "err", err, "policyName", nbResource.Spec.PolicyName)
return err
}
}
}
}
err := r.Client.Get(ctx, types.NamespacedName{Name: nbResource.Spec.PolicyName}, &nbPolicy)
if err != nil {
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", nbResource.Spec.PolicyName)
@@ -230,8 +251,6 @@ func (r *NBResourceReconciler) handleNetBirdResource(ctx context.Context, nbReso
return nil, err
}
nbResource.Status.NetworkResourceID = &resource.Id
} else if nbResource.Status.NetworkResourceID == nil && resource != nil {
nbResource.Status.NetworkResourceID = &resource.Id
} else if resource == nil {
// Status remembers networkResourceID but resource was deleted elsewhere
@@ -265,6 +284,44 @@ func (r *NBResourceReconciler) handleNetBirdResource(ctx context.Context, nbReso
// handleGroups create NBGroup objects for each group specified in NBResource
func (r *NBResourceReconciler) handleGroups(ctx context.Context, req ctrl.Request, nbResource *netbirdiov1.NBResource, logger logr.Logger) ([]string, *ctrl.Result, error) {
nbGroupList := netbirdiov1.NBGroupList{}
err := r.Client.List(ctx, &nbGroupList, &client.ListOptions{Namespace: req.Namespace})
if err != nil {
logger.Error(errKubernetesAPI, "error listing NBGroup", "err", err)
return nil, nil, err
}
for _, g := range nbGroupList.Items {
ownerIndex := -1
for idx, o := range g.OwnerReferences {
if o.UID == nbResource.UID {
ownerIndex = idx
break
}
}
if ownerIndex == -1 {
continue
}
if util.Contains(nbResource.Spec.Groups, g.Spec.Name) {
continue
}
if len(g.OwnerReferences) > 1 {
g.OwnerReferences = slices.Delete(g.OwnerReferences, ownerIndex, ownerIndex+1)
err = r.Client.Update(ctx, &g)
if err != nil && !errors.IsNotFound(err) {
logger.Error(errKubernetesAPI, "error updating NBGroup", "err", err)
return nil, nil, err
}
} else if len(g.OwnerReferences) == 1 {
g.Finalizers = util.Without(g.Finalizers, "netbird.io/resource-cleanup")
err = r.Client.Update(ctx, &g)
if err != nil && !errors.IsNotFound(err) {
logger.Error(errKubernetesAPI, "error updating NBGroup", "err", err)
return nil, nil, err
}
}
}
var groupIDs []string
for _, groupName := range nbResource.Spec.Groups {
@@ -283,8 +340,8 @@ func (r *NBResourceReconciler) handleGroups(ctx context.Context, req ctrl.Reques
Namespace: nbResource.Namespace,
OwnerReferences: []v1.OwnerReference{
{
APIVersion: nbResource.APIVersion,
Kind: nbResource.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: nbResource.Name,
UID: nbResource.UID,
BlockOwnerDeletion: util.Ptr(true),
@@ -315,8 +372,8 @@ func (r *NBResourceReconciler) handleGroups(ctx context.Context, req ctrl.Reques
if !ownerExists {
nbGroup.OwnerReferences = append(nbGroup.OwnerReferences, v1.OwnerReference{
APIVersion: nbResource.APIVersion,
Kind: nbResource.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: nbResource.Name,
UID: nbResource.UID,
BlockOwnerDeletion: util.Ptr(true),
@@ -380,8 +437,24 @@ func (r *NBResourceReconciler) handleDelete(ctx context.Context, req ctrl.Reques
}
for _, g := range nbGroupList.Items {
// TODO: Handle multiple owners
if len(g.OwnerReferences) > 0 && g.OwnerReferences[0].UID == nbResource.UID {
ownerIndex := -1
for idx, o := range g.OwnerReferences {
if o.UID == nbResource.UID {
ownerIndex = idx
break
}
}
if ownerIndex == -1 {
continue
}
if len(g.OwnerReferences) > 1 {
g.OwnerReferences = slices.Delete(g.OwnerReferences, ownerIndex, ownerIndex+1)
err = r.Client.Update(ctx, &g)
if err != nil && !errors.IsNotFound(err) {
logger.Error(errKubernetesAPI, "error updating NBGroup", "err", err)
return err
}
} else if len(g.OwnerReferences) == 1 {
g.Finalizers = util.Without(g.Finalizers, "netbird.io/resource-cleanup")
err = r.Client.Update(ctx, &g)
if err != nil && !errors.IsNotFound(err) {
+659 -22
View File
@@ -2,16 +2,25 @@ package controller
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"github.com/go-logr/logr"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
netbirdiov1 "github.com/netbirdio/kubernetes-operator/api/v1"
"github.com/netbirdio/kubernetes-operator/internal/util"
netbird "github.com/netbirdio/netbird/management/client/rest"
"github.com/netbirdio/netbird/management/server/http/api"
)
var _ = Describe("NBResource Controller", func() {
@@ -22,48 +31,676 @@ var _ = Describe("NBResource Controller", func() {
typeNamespacedName := types.NamespacedName{
Name: resourceName,
Namespace: "default", // TODO(user):Modify as needed
Namespace: "default",
}
nbresource := &netbirdiov1.NBResource{}
var netbirdClient *netbird.Client
var mux *http.ServeMux
var server *httptest.Server
var controllerReconciler *NBResourceReconciler
BeforeEach(func() {
Skip("Not implemented yet")
ctrl.SetLogger(logr.New(GinkgoLogr.GetSink()))
mux = &http.ServeMux{}
server = httptest.NewServer(mux)
netbirdClient = netbird.New(server.URL, "ABC")
controllerReconciler = &NBResourceReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
netbird: netbirdClient,
}
By("creating the custom resource for the Kind NBResource")
err := k8sClient.Get(ctx, typeNamespacedName, nbresource)
if err != nil && errors.IsNotFound(err) {
resource := &netbirdiov1.NBResource{
nbresource = &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: resourceName,
Namespace: "default",
Name: resourceName,
Namespace: "default",
Finalizers: []string{"netbird.io/cleanup"},
},
Spec: netbirdiov1.NBResourceSpec{
Name: "Test",
NetworkID: "test",
Address: "test.default.svc.cluster.local",
Groups: []string{"meow"},
TCPPorts: []int32{80},
},
// TODO(user): Specify other spec details if needed.
}
Expect(k8sClient.Create(ctx, resource)).To(Succeed())
Expect(k8sClient.Create(ctx, nbresource)).To(Succeed())
}
})
AfterEach(func() {
// TODO(user): Cleanup logic after each test, like removing the resource instance.
resource := &netbirdiov1.NBResource{}
err := k8sClient.Get(ctx, typeNamespacedName, resource)
Expect(err).NotTo(HaveOccurred())
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
By("Cleanup the specific resource instance NBResource")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
})
It("should successfully reconcile the resource", func() {
By("Reconciling the created resource")
controllerReconciler := &NBResourceReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
if len(resource.Finalizers) > 0 {
resource.Finalizers = nil
Expect(k8sClient.Update(ctx, resource)).To(Succeed())
}
By("Cleanup the specific resource instance NBResource")
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
BeforeEach(func() {
mux.HandleFunc("/api/groups", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
resp := []api.Group{
{
Id: "test",
Name: "meow",
},
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
})
})
When("Network Resource doesn't exist", Ordered, func() {
AfterAll(func() {
nbGroup := &netbirdiov1.NBGroup{}
err := k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)
if !errors.IsNotFound(err) {
if len(nbGroup.Finalizers) > 0 {
nbGroup.Finalizers = nil
Expect(k8sClient.Update(ctx, nbGroup)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, nbGroup)).To(Succeed())
}
})
It("should create NBGroups", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbGroup := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)).To(Succeed())
nbGroup.Status.GroupID = util.Ptr("test")
Expect(k8sClient.Status().Update(ctx, nbGroup)).To(Succeed())
})
It("should create Network Resource", func() {
networkResourceCreated := false
mux.HandleFunc("/api/networks/test/resources", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodPost {
networkResourceCreated = true
bs, err := io.ReadAll(r.Body)
Expect(err).NotTo(HaveOccurred())
var req api.PostApiNetworksNetworkIdResourcesJSONRequestBody
Expect(json.Unmarshal(bs, &req)).To(Succeed())
Expect(req.Name).To(Equal("Test"))
Expect(req.Description).NotTo(BeNil())
Expect(*req.Description).To(BeEquivalentTo("Created by kubernetes-operator"))
Expect(req.Enabled).To(BeTrue())
Expect(req.Groups).To(ConsistOf([]string{"test"}))
Expect(req.Address).To(Equal(nbresource.Spec.Address))
resp := api.NetworkResource{
Address: req.Address,
Description: req.Description,
Enabled: req.Enabled,
Groups: []api.GroupMinimum{
{
Id: "test",
Name: "meow",
},
},
Id: "test",
Name: req.Name,
Type: api.NetworkResourceTypeDomain,
}
bs, err = json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(networkResourceCreated).To(BeTrue())
})
})
When("Network Resource exists", func() {
BeforeEach(func() {
nbresource.Status.NetworkResourceID = util.Ptr("test")
Expect(k8sClient.Status().Update(ctx, nbresource)).To(Succeed())
nbGroup := &netbirdiov1.NBGroup{
ObjectMeta: metav1.ObjectMeta{
Name: "meow",
Namespace: "default",
Finalizers: []string{"netbird.io/resource-cleanup"},
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: "test-resource",
UID: nbresource.UID,
},
},
},
Spec: netbirdiov1.NBGroupSpec{
Name: "meow",
},
}
Expect(k8sClient.Create(ctx, nbGroup)).To(Succeed())
nbGroup.Status.GroupID = util.Ptr("test")
Expect(k8sClient.Status().Update(ctx, nbGroup)).To(Succeed())
})
AfterEach(func() {
nbGroup := &netbirdiov1.NBGroup{}
err := k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)
if !errors.IsNotFound(err) {
if len(nbGroup.Finalizers) > 0 {
nbGroup.Finalizers = nil
Expect(k8sClient.Update(ctx, nbGroup)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, nbGroup)).To(Succeed())
}
})
When("Network Resource is out of date", func() {
It("should update Network Resource", func() {
resourceUpdated := false
mux.HandleFunc("/api/networks/test/resources/test", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodGet {
resp := api.NetworkResource{
Address: nbresource.Spec.Address,
Description: &networkDescription,
Enabled: false,
Groups: []api.GroupMinimum{
{
Id: "test",
Name: "meow",
},
{
Id: "test2",
Name: "meow2",
},
},
Id: "test",
Name: nbresource.Spec.Name,
Type: api.NetworkResourceTypeDomain,
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
} else if r.Method == http.MethodPut {
resourceUpdated = true
bs, err := io.ReadAll(r.Body)
Expect(err).NotTo(HaveOccurred())
var req api.PutApiNetworksNetworkIdResourcesResourceIdJSONRequestBody
Expect(json.Unmarshal(bs, &req)).To(Succeed())
Expect(req.Name).To(Equal("Test"))
Expect(req.Description).NotTo(BeNil())
Expect(*req.Description).To(BeEquivalentTo("Created by kubernetes-operator"))
Expect(req.Enabled).To(BeTrue())
Expect(req.Groups).To(ConsistOf([]string{"test"}))
Expect(req.Address).To(Equal(nbresource.Spec.Address))
resp := api.NetworkResource{
Address: nbresource.Spec.Address,
Description: &networkDescription,
Enabled: true,
Groups: []api.GroupMinimum{
{
Id: "test",
Name: "meow",
},
},
Id: "test",
Name: nbresource.Spec.Name,
Type: api.NetworkResourceTypeDomain,
}
bs, err = json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
nbresource.Status.NetworkResourceID = util.Ptr("test")
Expect(k8sClient.Status().Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(resourceUpdated).To(BeTrue())
})
})
When("Network Resource is up-to-date", func() {
BeforeEach(func() {
mux.HandleFunc("/api/networks/test/resources/test", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodGet {
resp := api.NetworkResource{
Address: nbresource.Spec.Address,
Description: &networkDescription,
Enabled: true,
Groups: []api.GroupMinimum{
{
Id: "test",
Name: "meow",
},
},
Id: "test",
Name: nbresource.Spec.Name,
Type: api.NetworkResourceTypeDomain,
}
bs, err := json.Marshal(resp)
Expect(err).NotTo(HaveOccurred())
_, err = w.Write(bs)
Expect(err).NotTo(HaveOccurred())
}
})
})
When("Policy is specified", Ordered, func() {
BeforeAll(func() {
nbPolicy := &netbirdiov1.NBPolicy{
ObjectMeta: metav1.ObjectMeta{
Name: "test-a",
},
Spec: netbirdiov1.NBPolicySpec{
Name: "Test A",
SourceGroups: []string{"All"},
},
}
Expect(k8sClient.Create(ctx, nbPolicy)).To(Succeed())
nbPolicy = &netbirdiov1.NBPolicy{
ObjectMeta: metav1.ObjectMeta{
Name: "test-b",
},
Spec: netbirdiov1.NBPolicySpec{
Name: "Test B",
SourceGroups: []string{"All"},
},
}
Expect(k8sClient.Create(ctx, nbPolicy)).To(Succeed())
})
AfterAll(func() {
nbPolicy := &netbirdiov1.NBPolicy{}
err := k8sClient.Get(ctx, types.NamespacedName{Name: "test-a"}, nbPolicy)
if !errors.IsNotFound(err) {
Expect(k8sClient.Delete(ctx, nbPolicy)).To(Succeed())
}
nbPolicy = &netbirdiov1.NBPolicy{}
err = k8sClient.Get(ctx, types.NamespacedName{Name: "test-b"}, nbPolicy)
if !errors.IsNotFound(err) {
Expect(k8sClient.Delete(ctx, nbPolicy)).To(Succeed())
}
})
It("should update policy status", func() {
nbresource.Spec.PolicyName = "test-a"
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbPolicy := &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test-a"}, nbPolicy)).To(Succeed())
Expect(nbPolicy.Status.ManagedServiceList).To(ContainElement("default/test-resource"))
})
When("Policy is updated", func() {
It("should remove old reference and add new reference", func() {
nbresource.Spec.PolicyName = "test-b"
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
nbresource.Status.PolicyName = util.Ptr("test-a")
Expect(k8sClient.Status().Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbPolicy := &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test-a"}, nbPolicy)).To(Succeed())
Expect(nbPolicy.Status.ManagedServiceList).NotTo(ContainElement("default/test-resource"))
nbPolicy = &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test-b"}, nbPolicy)).To(Succeed())
Expect(nbPolicy.Status.ManagedServiceList).To(ContainElement("default/test-resource"))
})
})
When("Policy is removed", func() {
It("should remove old reference", func() {
nbPolicy := &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test-a"}, nbPolicy)).To(Succeed())
nbPolicy.Status.ManagedServiceList = []string{"default/test-resource"}
Expect(k8sClient.Status().Update(ctx, nbPolicy)).To(Succeed())
nbresource.Spec.PolicyName = ""
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
nbresource.Status.PolicyName = util.Ptr("test-a")
Expect(k8sClient.Status().Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbPolicy = &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test-a"}, nbPolicy)).To(Succeed())
Expect(nbPolicy.Status.ManagedServiceList).NotTo(ContainElement("default/test-resource"))
})
})
})
When("Groups are changed", func() {
When("Removed groups are no longer referenced by anything", func() {
It("should only remove finalizer", func() {
nbresource.Spec.Groups = []string{"meow2"}
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbGroup := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)).To(Succeed())
Expect(nbGroup.Finalizers).To(BeEmpty())
Expect(nbGroup.OwnerReferences).To(HaveLen(1))
})
})
When("Removed groups are referenced by something else", func() {
BeforeEach(func() {
otherResource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "not-test",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "nottest",
NetworkID: "test",
Address: "test",
Groups: []string{"test"},
},
}
Expect(k8sClient.Create(ctx, otherResource)).To(Succeed())
nbGroup := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)).To(Succeed())
nbGroup.OwnerReferences = append(nbGroup.OwnerReferences, metav1.OwnerReference{
APIVersion: "netbird.io/v1",
Kind: "NBResource",
Name: "not-test",
UID: otherResource.UID,
})
Expect(k8sClient.Update(ctx, nbGroup)).To(Succeed())
})
AfterEach(func() {
otherResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "not-test"}, otherResource)).To(Succeed())
Expect(k8sClient.Delete(ctx, otherResource)).To(Succeed())
})
It("should only remove owner reference", func() {
nbresource.Spec.Groups = []string{"meow2"}
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbGroup := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, nbGroup)).To(Succeed())
Expect(nbGroup.Finalizers).To(HaveLen(1))
Expect(nbGroup.OwnerReferences).To(HaveLen(1))
Expect(nbGroup.OwnerReferences[0].Name).To(Equal("not-test"))
})
})
When("New groups are added", func() {
It("should create new groups", func() {
nbresource.Spec.Groups = []string{"meow", "meow3"}
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbGroup := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow3"}, nbGroup)).To(Succeed())
Expect(nbGroup.OwnerReferences).To(HaveLen(1))
Expect(nbGroup.Finalizers).To(ConsistOf([]string{"netbird.io/group-cleanup", "netbird.io/resource-cleanup"}))
})
})
})
})
When("Network Resource is removed from NetBird", func() {
BeforeEach(func() {
mux.HandleFunc("/api/networks/test/resources/test", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodGet {
w.WriteHeader(404)
_, err := w.Write([]byte(`{"message": "not found", "code": 404}`))
Expect(err).NotTo(HaveOccurred())
}
})
})
It("should remove network resource ID and requeue", func() {
res, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(res.Requeue).To(BeTrue())
Expect(k8sClient.Get(ctx, typeNamespacedName, nbresource)).To(Succeed())
Expect(nbresource.Status.NetworkResourceID).To(BeNil())
})
})
})
When("NBResource is set for deletion", Ordered, func() {
BeforeAll(func() {
nbresource.Spec.Groups = []string{"meow", "meowdelete"}
Expect(k8sClient.Update(ctx, nbresource)).To(Succeed())
nbresource.Status.Groups = []string{"test", "testdelete"}
nbresource.Status.PolicyName = util.Ptr("test")
nbresource.Status.NetworkResourceID = util.Ptr("test")
Expect(k8sClient.Status().Update(ctx, nbresource)).To(Succeed())
nbPolicy := &netbirdiov1.NBPolicy{
ObjectMeta: metav1.ObjectMeta{
Name: "test",
},
Spec: netbirdiov1.NBPolicySpec{
Name: "Test",
SourceGroups: []string{"All"},
Bidirectional: true,
},
}
Expect(k8sClient.Create(ctx, nbPolicy)).To(Succeed())
nbPolicy.Status.ManagedServiceList = []string{"default/test-resource"}
Expect(k8sClient.Status().Update(ctx, nbPolicy)).To(Succeed())
nbGroup := &netbirdiov1.NBGroup{
ObjectMeta: metav1.ObjectMeta{
Name: "meow",
Namespace: "default",
Finalizers: []string{"netbird.io/resource-cleanup"},
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: nbresource.Name,
UID: nbresource.UID,
},
},
},
Spec: netbirdiov1.NBGroupSpec{
Name: "meow",
},
}
Expect(k8sClient.Create(ctx, nbGroup)).To(Succeed())
othernbresource := &netbirdiov1.NBResource{
ObjectMeta: metav1.ObjectMeta{
Name: "other-resource",
Namespace: "default",
},
Spec: netbirdiov1.NBResourceSpec{
Name: "test",
NetworkID: "test",
Address: "test",
Groups: []string{"test"},
},
}
Expect(k8sClient.Create(ctx, othernbresource)).To(Succeed())
nbGroup = &netbirdiov1.NBGroup{
ObjectMeta: metav1.ObjectMeta{
Name: "meowdelete",
Namespace: "default",
Finalizers: []string{"netbird.io/resource-cleanup"},
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: nbresource.Name,
UID: nbresource.UID,
},
{
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBResource",
Name: othernbresource.Name,
UID: othernbresource.UID,
},
},
},
Spec: netbirdiov1.NBGroupSpec{
Name: "meow",
},
}
Expect(k8sClient.Create(ctx, nbGroup)).To(Succeed())
})
AfterAll(func() {
policy := &netbirdiov1.NBPolicy{}
err := k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "test"}, policy)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
if len(policy.Finalizers) > 0 {
policy.Finalizers = nil
Expect(k8sClient.Update(ctx, policy)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, policy)).To(Succeed())
}
resource := &netbirdiov1.NBResource{}
err = k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "other-resource"}, resource)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
if len(resource.Finalizers) > 0 {
resource.Finalizers = nil
Expect(k8sClient.Update(ctx, resource)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
}
group := &netbirdiov1.NBGroup{}
err = k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, group)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
if len(group.Finalizers) > 0 {
group.Finalizers = nil
Expect(k8sClient.Update(ctx, group)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, group)).To(Succeed())
}
group = &netbirdiov1.NBGroup{}
err = k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meowdelete"}, group)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
if len(group.Finalizers) > 0 {
group.Finalizers = nil
Expect(k8sClient.Update(ctx, group)).To(Succeed())
}
Expect(k8sClient.Delete(ctx, group)).To(Succeed())
}
})
It("should delete Network Resource", func() {
Expect(k8sClient.Delete(ctx, nbresource)).To(Succeed())
resourceDeleted := false
mux.HandleFunc("/api/networks/test/resources/test", func(w http.ResponseWriter, r *http.Request) {
defer GinkgoRecover()
if r.Method == http.MethodDelete {
resourceDeleted = true
_, err := w.Write([]byte(`{}`))
Expect(err).NotTo(HaveOccurred())
}
})
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(resourceDeleted).To(BeTrue())
err = k8sClient.Get(ctx, typeNamespacedName, nbresource)
Expect(errors.IsNotFound(err)).To(BeTrue())
})
It("should remove resource cleanup finalizer from solely-owned NBGroups", func() {
group := &netbirdiov1.NBGroup{}
err := k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meow"}, group)
if errors.IsNotFound(err) {
return
}
Expect(group.Finalizers).To(BeEmpty())
})
It("should remove owner reference from shared NBGroups", func() {
group := &netbirdiov1.NBGroup{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "meowdelete"}, group)).To(Succeed())
Expect(group.Finalizers).To(HaveLen(1))
Expect(group.OwnerReferences).To(HaveLen(1))
})
It("should remove policy reference", func() {
policy := &netbirdiov1.NBPolicy{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Name: "test"}, policy)).To(Succeed())
Expect(policy.Status.ManagedServiceList).NotTo(ContainElement("default/test-resource"))
})
Expect(err).NotTo(HaveOccurred())
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
// Example: If you expect a certain status condition after reconciliation, verify it here.
})
})
})
+22 -13
View File
@@ -51,10 +51,10 @@ func (r *NBRoutingPeerReconciler) Reconcile(ctx context.Context, req ctrl.Reques
originalNBRP := nbrp.DeepCopy()
defer func() {
if originalNBRP.Status.NetworkID != nbrp.Status.NetworkID ||
originalNBRP.Status.RouterID != nbrp.Status.RouterID ||
originalNBRP.Status.SetupKeyID != nbrp.Status.SetupKeyID ||
!util.Equivalent(originalNBRP.Status.Conditions, nbrp.Status.Conditions) {
if originalNBRP.DeletionTimestamp != nil && len(nbrp.Finalizers) == 0 {
return
}
if !originalNBRP.Status.Equal(nbrp.Status) {
err = r.Client.Status().Update(ctx, nbrp)
if err != nil {
logger.Error(errKubernetesAPI, "error updating NBRoutingPeer Status", "err", err)
@@ -128,8 +128,8 @@ func (r *NBRoutingPeerReconciler) handleDeployment(ctx context.Context, req ctrl
Namespace: nbrp.Namespace,
OwnerReferences: []v1.OwnerReference{
{
APIVersion: nbrp.APIVersion,
Kind: nbrp.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBRoutingPeer",
Name: nbrp.Name,
UID: nbrp.UID,
BlockOwnerDeletion: util.Ptr(true),
@@ -199,8 +199,8 @@ func (r *NBRoutingPeerReconciler) handleDeployment(ctx context.Context, req ctrl
updatedDeployment.ObjectMeta.Namespace = nbrp.Namespace
updatedDeployment.ObjectMeta.OwnerReferences = []v1.OwnerReference{
{
APIVersion: nbrp.APIVersion,
Kind: nbrp.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBRoutingPeer",
Name: nbrp.Name,
UID: nbrp.UID,
BlockOwnerDeletion: util.Ptr(true),
@@ -355,8 +355,8 @@ func (r *NBRoutingPeerReconciler) handleSetupKey(ctx context.Context, req ctrl.R
Namespace: nbrp.Namespace,
OwnerReferences: []v1.OwnerReference{
{
APIVersion: nbrp.APIVersion,
Kind: nbrp.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBRoutingPeer",
Name: nbrp.Name,
UID: nbrp.UID,
BlockOwnerDeletion: util.Ptr(true),
@@ -387,8 +387,17 @@ func (r *NBRoutingPeerReconciler) handleSetupKey(ctx context.Context, req ctrl.R
}
if (err != nil && strings.Contains(err.Error(), "not found")) || setupKey.Revoked {
nbrp.Status.SetupKeyID = nil
if setupKey != nil && setupKey.Revoked {
err = r.netbird.SetupKeys.Delete(ctx, *nbrp.Status.SetupKeyID)
if err != nil {
logger.Error(errNetBirdAPI, "error deleting setup key", "err", err)
nbrp.Status.Conditions = netbirdiov1.NBConditionFalse("APIError", fmt.Sprintf("error deleting setup key: %v", err))
return &ctrl.Result{}, err
}
}
nbrp.Status.SetupKeyID = nil
// Requeue to avoid repeating code
return &ctrl.Result{Requeue: true}, nil
}
@@ -447,8 +456,8 @@ func (r *NBRoutingPeerReconciler) handleGroup(ctx context.Context, req ctrl.Requ
Namespace: nbrp.Namespace,
OwnerReferences: []v1.OwnerReference{
{
APIVersion: nbrp.APIVersion,
Kind: nbrp.Kind,
APIVersion: netbirdiov1.GroupVersion.Identifier(),
Kind: "NBRoutingPeer",
Name: nbrp.Name,
UID: nbrp.UID,
BlockOwnerDeletion: util.Ptr(true),
File diff suppressed because it is too large Load Diff
+57 -41
View File
@@ -157,7 +157,8 @@ func (r *ServiceReconciler) exposeService(ctx context.Context, req ctrl.Request,
return ctrl.Result{}, err
}
nbrsErr := r.reconcileNBResource(&nbResource, req, svc, routingPeer)
originalNBResource := nbResource.DeepCopy()
nbrsErr := r.reconcileNBResource(&nbResource, req, svc, routingPeer, logger)
if nbrsErr != nil {
return ctrl.Result{}, nbrsErr
}
@@ -168,7 +169,7 @@ func (r *ServiceReconciler) exposeService(ctx context.Context, req ctrl.Request,
logger.Error(errKubernetesAPI, "error creating NBResource", "err", err)
return ctrl.Result{}, err
}
} else {
} else if !originalNBResource.Spec.Equal(nbResource.Spec) {
err = r.Client.Update(ctx, &nbResource)
if err != nil {
logger.Error(errKubernetesAPI, "error updating NBResource", "err", err)
@@ -180,7 +181,7 @@ func (r *ServiceReconciler) exposeService(ctx context.Context, req ctrl.Request,
}
// reconcileNBResource ensures NBResource settings are in-line with Service definition and annotations
func (r *ServiceReconciler) reconcileNBResource(nbResource *netbirdiov1.NBResource, req ctrl.Request, svc corev1.Service, routingPeer netbirdiov1.NBRoutingPeer) error {
func (r *ServiceReconciler) reconcileNBResource(nbResource *netbirdiov1.NBResource, req ctrl.Request, svc corev1.Service, routingPeer netbirdiov1.NBRoutingPeer, logger logr.Logger) error {
groups := []string{fmt.Sprintf("%s-%s-%s", r.ClusterName, req.Namespace, req.Name)}
if v, ok := svc.Annotations[serviceGroupsAnnotation]; ok {
groups = nil
@@ -202,50 +203,65 @@ func (r *ServiceReconciler) reconcileNBResource(nbResource *netbirdiov1.NBResour
nbResource.Spec.Address = fmt.Sprintf("%s.%s.%s", svc.Name, svc.Namespace, r.ClusterDNS)
nbResource.Spec.Groups = groups
if v, ok := svc.Annotations[servicePolicyAnnotation]; ok {
nbResource.Spec.PolicyName = v
var filterProtocols []string
if v, ok := svc.Annotations[serviceProtocolAnnotation]; ok {
filterProtocols = []string{v}
if _, ok := svc.Annotations[servicePolicyAnnotation]; ok {
err := r.applyPolicy(nbResource, svc, logger)
if err != nil {
return err
}
var filterPorts []int32
if v, ok := svc.Annotations[servicePortsAnnotation]; ok {
for _, v := range strings.Split(v, ",") {
port, err := strconv.ParseInt(v, 10, 64)
if err != nil {
return err
}
} else if nbResource.Spec.PolicyName != "" {
nbResource.Spec.PolicyName = ""
nbResource.Spec.TCPPorts = nil
nbResource.Spec.UDPPorts = nil
}
filterPorts = append(filterPorts, int32(port))
}
}
return nil
}
for _, p := range svc.Spec.Ports {
if len(filterProtocols) > 0 && !util.Contains(filterProtocols, string(p.Protocol)) {
continue
}
if len(filterPorts) > 0 && !util.Contains(filterPorts, p.Port) {
continue
}
switch p.Protocol {
case corev1.ProtocolSCTP:
if !util.Contains(nbResource.Spec.TCPPorts, p.Port) {
nbResource.Spec.TCPPorts = append(nbResource.Spec.TCPPorts, p.Port)
}
case corev1.ProtocolTCP:
if !util.Contains(nbResource.Spec.TCPPorts, p.Port) {
nbResource.Spec.TCPPorts = append(nbResource.Spec.TCPPorts, p.Port)
}
case corev1.ProtocolUDP:
if !util.Contains(nbResource.Spec.UDPPorts, p.Port) {
nbResource.Spec.UDPPorts = append(nbResource.Spec.UDPPorts, p.Port)
}
default:
return errUnknownProtocol
func (r *ServiceReconciler) applyPolicy(nbResource *netbirdiov1.NBResource, svc corev1.Service, logger logr.Logger) error {
nbResource.Spec.PolicyName = svc.Annotations[servicePolicyAnnotation]
var filterProtocols []string
if v, ok := svc.Annotations[serviceProtocolAnnotation]; ok {
filterProtocols = []string{v}
}
var filterPorts []int32
if v, ok := svc.Annotations[servicePortsAnnotation]; ok {
for _, v := range strings.Split(v, ",") {
port, err := strconv.ParseInt(v, 10, 64)
if err != nil {
return err
}
filterPorts = append(filterPorts, int32(port))
}
}
for _, p := range svc.Spec.Ports {
switch p.Protocol {
case corev1.ProtocolTCP:
if (len(filterPorts) > 0 && !util.Contains(filterPorts, p.Port)) || (len(filterProtocols) > 0 && !util.Contains(filterProtocols, "tcp")) {
if util.Contains(nbResource.Spec.TCPPorts, p.Port) {
nbResource.Spec.TCPPorts = util.Without(nbResource.Spec.TCPPorts, p.Port)
}
continue
}
if !util.Contains(nbResource.Spec.TCPPorts, p.Port) {
nbResource.Spec.TCPPorts = append(nbResource.Spec.TCPPorts, p.Port)
}
case corev1.ProtocolUDP:
if (len(filterPorts) > 0 && !util.Contains(filterPorts, p.Port)) || (len(filterProtocols) > 0 && !util.Contains(filterProtocols, "udp")) {
if util.Contains(nbResource.Spec.UDPPorts, p.Port) {
nbResource.Spec.UDPPorts = util.Without(nbResource.Spec.UDPPorts, p.Port)
}
continue
}
if !util.Contains(nbResource.Spec.UDPPorts, p.Port) {
nbResource.Spec.UDPPorts = append(nbResource.Spec.UDPPorts, p.Port)
}
default:
logger.Info("Unsupported protocol %v", p.Protocol)
continue
}
}
// TODO: Handle removed policy name
return nil
}
+460 -4
View File
@@ -1,17 +1,473 @@
package controller
import (
netbirdiov1 "github.com/netbirdio/kubernetes-operator/api/v1"
"github.com/netbirdio/kubernetes-operator/internal/util"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/intstr"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)
var _ = Describe("Service Controller", func() {
Context("When reconciling a resource", func() {
typeNamespacedName := types.NamespacedName{
Namespace: "default",
Name: "test-resource",
}
const policyName = "test"
var service *corev1.Service
It("should successfully reconcile the resource", func() {
Skip("Not implemented yet")
var controllerReconciler *ServiceReconciler
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
// Example: If you expect a certain status condition after reconciliation, verify it here.
BeforeEach(func() {
service = &corev1.Service{
ObjectMeta: v1.ObjectMeta{
Name: "test-resource",
Namespace: "default",
},
Spec: corev1.ServiceSpec{
Ports: []corev1.ServicePort{
{
Name: "a",
Protocol: corev1.ProtocolTCP,
Port: 80,
TargetPort: intstr.FromInt(80),
},
{
Name: "b",
Protocol: corev1.ProtocolTCP,
Port: 443,
TargetPort: intstr.FromInt(443),
},
{
Name: "c",
Protocol: corev1.ProtocolUDP,
Port: 80,
TargetPort: intstr.FromInt(80),
},
{
Name: "d",
Protocol: corev1.ProtocolUDP,
Port: 443,
TargetPort: intstr.FromInt(443),
},
},
},
}
Expect(k8sClient.Create(ctx, service)).To(Succeed())
controllerReconciler = &ServiceReconciler{
Client: k8sClient,
Scheme: k8sClient.Scheme(),
ClusterName: "kubernetes",
NamespacedNetworks: false,
ClusterDNS: "svc.cluster.local",
ControllerNamespace: "default",
}
})
AfterEach(func() {
svc := &corev1.Service{}
err := k8sClient.Get(ctx, typeNamespacedName, svc)
if !errors.IsNotFound(err) {
if len(svc.Finalizers) > 0 {
svc.Finalizers = nil
Expect(k8sClient.Update(ctx, svc)).To(Succeed())
}
err := k8sClient.Delete(ctx, svc)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
}
}
nbrp := &netbirdiov1.NBRoutingPeer{}
err = k8sClient.Get(ctx, types.NamespacedName{Namespace: "default", Name: "router"}, nbrp)
if !errors.IsNotFound(err) {
if len(nbrp.Finalizers) > 0 {
nbrp.Finalizers = nil
Expect(k8sClient.Update(ctx, nbrp)).To(Succeed())
}
err := k8sClient.Delete(ctx, nbrp)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
}
}
nbResource := &netbirdiov1.NBResource{}
err = k8sClient.Get(ctx, typeNamespacedName, nbResource)
if !errors.IsNotFound(err) {
if len(nbResource.Finalizers) > 0 {
nbResource.Finalizers = nil
Expect(k8sClient.Update(ctx, nbResource)).To(Succeed())
}
err := k8sClient.Delete(ctx, nbResource)
if !errors.IsNotFound(err) {
Expect(err).NotTo(HaveOccurred())
}
}
})
When("Service is not already exposed", func() {
When("Service should not be exposed", func() {
It("should change nothing", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(k8sClient.Get(ctx, typeNamespacedName, service)).To(Succeed())
Expect(service.Finalizers).To(BeEmpty())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).NotTo(Succeed())
})
})
When("NBRoutingPeer doesn't exist", func() {
BeforeEach(func() {
if service.Annotations == nil {
service.Annotations = make(map[string]string)
}
service.Annotations[ServiceExposeAnnotation] = "trueish"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
})
It("should create NBRoutingPeer and requeue until network ID is available", func() {
res, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(res.RequeueAfter).NotTo(BeZero())
nbrp := &netbirdiov1.NBRoutingPeer{}
Expect(k8sClient.Get(ctx, types.NamespacedName{Namespace: typeNamespacedName.Namespace, Name: "router"}, nbrp)).To(Succeed())
res, err = controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(res.RequeueAfter).NotTo(BeZero())
nbrp.Status.NetworkID = util.Ptr(policyName)
Expect(k8sClient.Status().Update(ctx, nbrp)).To(Succeed())
res, err = controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(res.RequeueAfter).To(BeZero())
})
})
When("NBRoutingPeer exists", func() {
BeforeEach(func() {
nbrp := &netbirdiov1.NBRoutingPeer{
ObjectMeta: v1.ObjectMeta{
Namespace: typeNamespacedName.Namespace,
Name: "router",
},
Spec: netbirdiov1.NBRoutingPeerSpec{},
}
Expect(k8sClient.Create(ctx, nbrp)).To(Succeed())
nbrp.Status.NetworkID = util.Ptr(policyName)
Expect(k8sClient.Status().Update(ctx, nbrp)).To(Succeed())
})
When("Service should be exposed", func() {
BeforeEach(func() {
if service.Annotations == nil {
service.Annotations = make(map[string]string)
}
service.Annotations[ServiceExposeAnnotation] = "true"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
})
It("should add finalizer to service object", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(k8sClient.Get(ctx, typeNamespacedName, service)).To(Succeed())
Expect(service.Finalizers).To(ContainElement("netbird.io/cleanup"))
})
When("nothing else is specified", func() {
It("should create NBResource with default values", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.Address).To(Equal(typeNamespacedName.Name + "." + typeNamespacedName.Namespace + "." + controllerReconciler.ClusterDNS))
Expect(nbResource.Spec.Groups).To(ConsistOf([]string{controllerReconciler.ClusterName + "-" + typeNamespacedName.Namespace + "-" + typeNamespacedName.Name}))
Expect(nbResource.Spec.Name).To(Equal(typeNamespacedName.Namespace + "-" + typeNamespacedName.Name))
Expect(nbResource.Spec.NetworkID).To(Equal(policyName))
Expect(nbResource.Spec.PolicyName).To(BeEmpty())
Expect(nbResource.Spec.TCPPorts).To(BeEmpty())
Expect(nbResource.Spec.UDPPorts).To(BeEmpty())
})
})
When("policy is specified", func() {
BeforeEach(func() {
service.Annotations[servicePolicyAnnotation] = policyName
Expect(k8sClient.Update(ctx, service)).To(Succeed())
})
When("nothing is restricted", func() {
It("should create NBResource with policy", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.PolicyName).To(Equal(policyName))
Expect(nbResource.Spec.TCPPorts).To(ConsistOf([]int32{443, 80}))
Expect(nbResource.Spec.UDPPorts).To(ConsistOf([]int32{443, 80}))
})
})
When("ports are restricted", func() {
It("should create NBResource with policy and only specified ports", func() {
service.Annotations[servicePortsAnnotation] = "80"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.PolicyName).To(Equal(policyName))
Expect(nbResource.Spec.TCPPorts).To(ConsistOf([]int32{80}))
Expect(nbResource.Spec.UDPPorts).To(ConsistOf([]int32{80}))
})
})
When("protocol is restricted", func() {
It("should create NBResource with policy and only specified protocol", func() {
service.Annotations[serviceProtocolAnnotation] = "tcp"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.PolicyName).To(Equal(policyName))
Expect(nbResource.Spec.TCPPorts).To(ConsistOf([]int32{80, 443}))
Expect(nbResource.Spec.UDPPorts).To(BeEmpty())
})
})
})
When("resource name is specified", func() {
It("should create NBResource with specified name", func() {
service.Annotations[serviceResourceAnnotation] = "meow"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.Name).To(Equal("meow"))
})
})
When("resource groups specified", func() {
It("should create NBResource with specified groups", func() {
service.Annotations[serviceGroupsAnnotation] = "meow, wow ,test"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.Groups).To(ConsistOf([]string{"meow", "wow", policyName}))
})
})
})
})
})
When("Service is already exposed", func() {
BeforeEach(func() {
nbResource := &netbirdiov1.NBResource{
ObjectMeta: v1.ObjectMeta{
Name: typeNamespacedName.Name,
Namespace: typeNamespacedName.Namespace,
},
Spec: netbirdiov1.NBResourceSpec{
Name: typeNamespacedName.Namespace + "-" + typeNamespacedName.Name,
Address: typeNamespacedName.Name + "." + typeNamespacedName.Namespace + "." + controllerReconciler.ClusterDNS,
Groups: []string{controllerReconciler.ClusterName + "-" + typeNamespacedName.Namespace + "-" + typeNamespacedName.Name},
NetworkID: policyName,
},
}
Expect(k8sClient.Create(ctx, nbResource)).To(Succeed())
if service.Annotations == nil {
service.Annotations = make(map[string]string)
}
service.Annotations[ServiceExposeAnnotation] = "true"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
nbrp := &netbirdiov1.NBRoutingPeer{
ObjectMeta: v1.ObjectMeta{
Namespace: typeNamespacedName.Namespace,
Name: "router",
},
Spec: netbirdiov1.NBRoutingPeerSpec{},
}
Expect(k8sClient.Create(ctx, nbrp)).To(Succeed())
nbrp.Status.NetworkID = util.Ptr(policyName)
Expect(k8sClient.Status().Update(ctx, nbrp)).To(Succeed())
})
When("Service should not be exposed", func() {
BeforeEach(func() {
delete(service.Annotations, ServiceExposeAnnotation)
Expect(k8sClient.Update(ctx, service)).To(Succeed())
})
It("should delete NBResource", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
err = k8sClient.Get(ctx, typeNamespacedName, nbResource)
if !errors.IsNotFound(err) {
Expect(nbResource.DeletionTimestamp).NotTo(BeNil())
}
})
It("should remove finalizer from Service", func() {
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
Expect(k8sClient.Get(ctx, typeNamespacedName, service)).To(Succeed())
Expect(service.Finalizers).NotTo(ContainElement("netbird.io/cleanup"))
})
})
When("Nothing changes", func() {
It("should do nothing", func() {
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
resourceVersion := nbResource.ResourceVersion
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource = &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(resourceVersion).To(BeEquivalentTo(nbResource.ResourceVersion))
})
})
When("policy changes", func() {
It("should update policy in NBResource spec", func() {
service.Annotations[servicePolicyAnnotation] = policyName
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.PolicyName).To(Equal(policyName))
})
})
When("policy is removed", func() {
It("should remove policy in NBResource spec", func() {
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
nbResource.Spec.PolicyName = policyName
Expect(k8sClient.Update(ctx, nbResource)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource = &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.PolicyName).To(Equal(""))
})
})
When("policy ports changes", func() {
It("should update ports in NBResource spec", func() {
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
nbResource.Spec.PolicyName = policyName
nbResource.Spec.TCPPorts = []int32{443, 80}
nbResource.Spec.UDPPorts = []int32{443, 80}
Expect(k8sClient.Update(ctx, nbResource)).To(Succeed())
service.Annotations[servicePolicyAnnotation] = policyName
service.Annotations[servicePortsAnnotation] = "80"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource = &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.TCPPorts).To(ConsistOf([]int32{80}))
Expect(nbResource.Spec.UDPPorts).To(ConsistOf([]int32{80}))
})
})
When("policy protocol changes", func() {
It("should update protocol in NBResource spec", func() {
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
nbResource.Spec.PolicyName = policyName
nbResource.Spec.TCPPorts = []int32{443, 80}
nbResource.Spec.UDPPorts = []int32{443, 80}
Expect(k8sClient.Update(ctx, nbResource)).To(Succeed())
service.Annotations[servicePolicyAnnotation] = policyName
service.Annotations[serviceProtocolAnnotation] = "tcp"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource = &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.TCPPorts).To(ConsistOf([]int32{80, 443}))
Expect(nbResource.Spec.UDPPorts).To(BeEmpty())
})
})
When("resource name changes", func() {
It("should update name in NBResource spec", func() {
service.Annotations[serviceResourceAnnotation] = "meow"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.Name).To(Equal("meow"))
})
})
When("resource groups changes", func() {
It("should update groups in NBResource spec", func() {
service.Annotations[serviceGroupsAnnotation] = "a7medmo7sen, pewpewpew"
Expect(k8sClient.Update(ctx, service)).To(Succeed())
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
NamespacedName: typeNamespacedName,
})
Expect(err).NotTo(HaveOccurred())
nbResource := &netbirdiov1.NBResource{}
Expect(k8sClient.Get(ctx, typeNamespacedName, nbResource)).To(Succeed())
Expect(nbResource.Spec.Groups).To(ConsistOf([]string{"a7medmo7sen", "pewpewpew"}))
})
})
})
})
})