mirror of
https://github.com/YuzuZensai/netbird-kubernetes-operator.git
synced 2026-09-13 10:49:15 +00:00
Add support for policy auto-creation
This commit is contained in:
@@ -14,6 +14,7 @@ import (
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/handler"
|
||||
"sigs.k8s.io/controller-runtime/pkg/reconcile"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
netbirdiov1 "github.com/netbirdio/kubernetes-operator/api/v1"
|
||||
@@ -25,10 +26,12 @@ import (
|
||||
// NBResourceReconciler reconciles a NBResource object
|
||||
type NBResourceReconciler struct {
|
||||
client.Client
|
||||
Scheme *runtime.Scheme
|
||||
APIKey string
|
||||
ManagementURL string
|
||||
netbird *netbird.Client
|
||||
Scheme *runtime.Scheme
|
||||
APIKey string
|
||||
ManagementURL string
|
||||
AllowAutomaticPolicyCreation bool
|
||||
ClusterName string
|
||||
netbird *netbird.Client
|
||||
}
|
||||
|
||||
// Reconcile is part of the main kubernetes reconciliation loop which aims to
|
||||
@@ -96,6 +99,7 @@ func (r *NBResourceReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
err = r.handlePolicy(ctx, req, nbResource, groupIDs, logger)
|
||||
if err != nil {
|
||||
nbResource.Status.Conditions = netbirdiov1.NBConditionFalse("internalError", fmt.Sprintf("Error occurred handling policy changes: %v", err))
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
nbResource.Status.Conditions = netbirdiov1.NBConditionTrue()
|
||||
@@ -103,115 +107,236 @@ func (r *NBResourceReconciler) Reconcile(ctx context.Context, req ctrl.Request)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
func (r *NBResourceReconciler) handlePolicyCreate(ctx context.Context, nbResource *netbirdiov1.NBResource, req ctrl.Request, policy string, nbPolicy *netbirdiov1.NBPolicy, logger logr.Logger) error {
|
||||
if len(nbResource.Spec.PolicySourceGroups) == 0 {
|
||||
logger.Error(errInvalidValue, "Cannot auto-generate policy, missing source groups.")
|
||||
return fmt.Errorf("cannot auto-generate policy, missing source groups")
|
||||
}
|
||||
name := nbResource.Spec.PolicyFriendlyName[policy]
|
||||
if name == "" {
|
||||
name = fmt.Sprintf("Autogenerated policy for resource %s/%s in cluster %s", nbResource.Namespace, nbResource.Name, r.ClusterName)
|
||||
}
|
||||
generatedName := fmt.Sprintf("%s-%s-%s", policy, req.Namespace, req.Name)
|
||||
*nbPolicy = netbirdiov1.NBPolicy{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: generatedName,
|
||||
Annotations: map[string]string{"netbird.io/generated-by": req.NamespacedName.String()},
|
||||
Finalizers: []string{"netbird.io/cleanup"},
|
||||
},
|
||||
Spec: netbirdiov1.NBPolicySpec{
|
||||
Name: name,
|
||||
Description: "Generated by " + req.NamespacedName.String(),
|
||||
SourceGroups: nbResource.Spec.PolicySourceGroups,
|
||||
Bidirectional: true,
|
||||
},
|
||||
}
|
||||
|
||||
err := r.Client.Create(ctx, nbPolicy)
|
||||
if errors.IsAlreadyExists(err) {
|
||||
err = r.Client.Get(ctx, types.NamespacedName{Name: generatedName}, nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
if nbPolicy.Annotations == nil {
|
||||
nbPolicy.Annotations = make(map[string]string)
|
||||
}
|
||||
nbPolicy.Annotations["netbird.io/generated-by"] = req.NamespacedName.String()
|
||||
nbPolicy.Spec = netbirdiov1.NBPolicySpec{
|
||||
Name: name,
|
||||
Description: "Generated by " + req.NamespacedName.String(),
|
||||
SourceGroups: nbResource.Spec.PolicySourceGroups,
|
||||
Bidirectional: true,
|
||||
}
|
||||
|
||||
err = r.Client.Update(ctx, nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "err", err)
|
||||
return err
|
||||
}
|
||||
} else if err != nil {
|
||||
logger.Error(errKubernetesAPI, "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
if nbResource.Status.PolicyNameMapping == nil {
|
||||
nbResource.Status.PolicyNameMapping = make(map[string]string)
|
||||
}
|
||||
nbResource.Status.PolicyNameMapping[policy] = generatedName
|
||||
nbResource.Status.PolicySourceGroups = nbResource.Spec.PolicySourceGroups
|
||||
nbResource.Status.PolicyFriendlyName = nbResource.Spec.PolicyFriendlyName
|
||||
|
||||
nbPolicy.Status.ManagedServiceList = append(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
err = r.Client.Status().Update(ctx, nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "err", err)
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *NBResourceReconciler) handlePolicyAddUpdate(ctx context.Context, req ctrl.Request, nbResource *netbirdiov1.NBResource, policy string, groupIDs []string, logger logr.Logger) error {
|
||||
var nbPolicy netbirdiov1.NBPolicy
|
||||
updatePolicyStatus := false
|
||||
|
||||
kubernetesPolicyName := policy
|
||||
if v, ok := nbResource.Status.PolicyNameMapping[policy]; ok {
|
||||
kubernetesPolicyName = v
|
||||
}
|
||||
err := r.Client.Get(ctx, types.NamespacedName{Name: kubernetesPolicyName}, &nbPolicy)
|
||||
if errors.IsNotFound(err) && r.AllowAutomaticPolicyCreation {
|
||||
err = r.handlePolicyCreate(ctx, nbResource, req, policy, &nbPolicy, logger)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
} else if errors.IsNotFound(err) && !r.AllowAutomaticPolicyCreation {
|
||||
logger.Info("automatic policy creation is not allowed")
|
||||
return nil
|
||||
} else if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
|
||||
if !util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
|
||||
nbPolicy.Status.ManagedServiceList = append(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Spec.TCPPorts, nbResource.Status.TCPPorts) {
|
||||
nbResource.Status.TCPPorts = nbResource.Spec.TCPPorts
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Spec.UDPPorts, nbResource.Status.UDPPorts) {
|
||||
nbResource.Status.UDPPorts = nbResource.Spec.UDPPorts
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Status.Groups, groupIDs) {
|
||||
nbResource.Status.Groups = groupIDs
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if _, ok := nbResource.Status.PolicyNameMapping[policy]; ok {
|
||||
updatePolicySpec := false
|
||||
if v, ok := nbPolicy.Annotations["netbird.io/generated-by"]; !ok || v != req.NamespacedName.String() {
|
||||
if nbPolicy.Annotations == nil {
|
||||
nbPolicy.Annotations = make(map[string]string)
|
||||
}
|
||||
nbPolicy.Annotations["netbird.io/generated-by"] = req.NamespacedName.String()
|
||||
updatePolicySpec = true
|
||||
}
|
||||
|
||||
if v, ok := nbResource.Spec.PolicyFriendlyName[policy]; ok {
|
||||
if nbPolicy.Spec.Name != v {
|
||||
nbPolicy.Spec.Name = v
|
||||
updatePolicySpec = true
|
||||
}
|
||||
} else {
|
||||
if nbPolicy.Spec.Name != fmt.Sprintf("Autogenerated policy for resource %s/%s in cluster %s", nbResource.Namespace, nbResource.Name, r.ClusterName) {
|
||||
nbPolicy.Spec.Name = fmt.Sprintf("Autogenerated policy for resource %s/%s in cluster %s", nbResource.Namespace, nbResource.Name, r.ClusterName)
|
||||
updatePolicySpec = true
|
||||
}
|
||||
}
|
||||
|
||||
if nbPolicy.Spec.Description != "Generated by "+req.NamespacedName.String() {
|
||||
nbPolicy.Spec.Description = "Generated by " + req.NamespacedName.String()
|
||||
updatePolicySpec = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbPolicy.Spec.SourceGroups, nbResource.Spec.PolicySourceGroups) {
|
||||
nbPolicy.Spec.SourceGroups = nbResource.Spec.PolicySourceGroups
|
||||
updatePolicySpec = true
|
||||
}
|
||||
|
||||
if updatePolicySpec {
|
||||
err := r.Client.Update(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if updatePolicyStatus {
|
||||
err := r.Client.Status().Update(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error updating NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *NBResourceReconciler) handlePolicyDelete(ctx context.Context, req ctrl.Request, nbResource *netbirdiov1.NBResource, specPolicies []string, policy string, logger logr.Logger) error {
|
||||
var nbPolicy netbirdiov1.NBPolicy
|
||||
if !util.Contains(specPolicies, policy) {
|
||||
kubeName := policy
|
||||
if v, ok := nbResource.Status.PolicyNameMapping[policy]; ok {
|
||||
kubeName = v
|
||||
}
|
||||
err := r.Client.Get(ctx, types.NamespacedName{Name: kubeName}, &nbPolicy)
|
||||
if !errors.IsNotFound(err) {
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
|
||||
if _, ok := nbResource.Status.PolicyNameMapping[policy]; ok {
|
||||
// Delete Policy
|
||||
err := r.Client.Delete(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error deleting NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
|
||||
delete(nbResource.Status.PolicyNameMapping, policy)
|
||||
} else if util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
|
||||
nbPolicy.Status.ManagedServiceList = util.Without(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
err := r.Client.Status().Update(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error updating NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// handlePolicy update NBPolicy if defined to add self reference to policy status
|
||||
func (r *NBResourceReconciler) handlePolicy(ctx context.Context, req ctrl.Request, nbResource *netbirdiov1.NBResource, groupIDs []string, logger logr.Logger) error {
|
||||
if nbResource.Status.PolicyName == nil && nbResource.Spec.PolicyName == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
var nbPolicy netbirdiov1.NBPolicy
|
||||
if nbResource.Spec.PolicyName == "" && nbResource.Status.PolicyName != nil {
|
||||
// Remove self reference from policy status
|
||||
policies := util.SplitTrim(*nbResource.Status.PolicyName, ",")
|
||||
for _, policyName := range policies {
|
||||
err := r.Client.Get(ctx, types.NamespacedName{Name: policyName}, &nbPolicy)
|
||||
nbResource.Status.PolicyName = nil
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", policyName)
|
||||
return err
|
||||
}
|
||||
if util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
|
||||
nbPolicy.Status.ManagedServiceList = util.Without(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
err := r.Client.Status().Update(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error updating NBPolicy", "err", err, "policyName", policyName)
|
||||
return err
|
||||
}
|
||||
}
|
||||
specPolicies := util.SplitTrim(nbResource.Spec.PolicyName, ",")
|
||||
var statusPolicies []string
|
||||
if nbResource.Status.PolicyName != nil {
|
||||
statusPolicies = util.SplitTrim(*nbResource.Status.PolicyName, ",")
|
||||
}
|
||||
|
||||
for _, policy := range specPolicies {
|
||||
err := r.handlePolicyAddUpdate(ctx, req, nbResource, policy, groupIDs, logger)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
specPolicies := util.SplitTrim(nbResource.Spec.PolicyName, ",")
|
||||
var statusPolicies []string
|
||||
if nbResource.Status.PolicyName != nil {
|
||||
statusPolicies = util.SplitTrim(*nbResource.Status.PolicyName, ",")
|
||||
}
|
||||
|
||||
for _, policy := range statusPolicies {
|
||||
err := r.handlePolicyDelete(ctx, req, nbResource, specPolicies, policy, logger)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
for _, policy := range specPolicies {
|
||||
updatePolicyStatus := false
|
||||
|
||||
err := r.Client.Get(ctx, types.NamespacedName{Name: policy}, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
if !util.Contains(statusPolicies, policy) {
|
||||
// New
|
||||
if !util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
|
||||
nbPolicy.Status.ManagedServiceList = append(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
} else {
|
||||
// Check update
|
||||
if !util.Contains(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String()) {
|
||||
nbPolicy.Status.ManagedServiceList = append(nbPolicy.Status.ManagedServiceList, req.NamespacedName.String())
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Spec.TCPPorts, nbResource.Status.TCPPorts) {
|
||||
nbResource.Status.TCPPorts = nbResource.Spec.TCPPorts
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Spec.UDPPorts, nbResource.Status.UDPPorts) {
|
||||
nbResource.Status.UDPPorts = nbResource.Spec.UDPPorts
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
|
||||
if !util.Equivalent(nbResource.Status.Groups, groupIDs) {
|
||||
nbResource.Status.Groups = groupIDs
|
||||
nbPolicy.Status.LastUpdatedAt = &v1.Time{Time: time.Now()}
|
||||
updatePolicyStatus = true
|
||||
}
|
||||
}
|
||||
|
||||
if updatePolicyStatus {
|
||||
err := r.Client.Status().Update(ctx, &nbPolicy)
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error updating NBPolicy", "err", err, "policyName", policy)
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for _, policy := range statusPolicies {
|
||||
// Delete
|
||||
if !util.Contains(specPolicies, policy) {
|
||||
err := r.Client.Get(ctx, types.NamespacedName{Name: policy}, &nbPolicy)
|
||||
if !errors.IsNotFound(err) {
|
||||
if err != nil {
|
||||
logger.Error(errKubernetesAPI, "error getting NBPolicy", "err", err, "policyName", policy)
|
||||
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", policy)
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if nbResource.Status.PolicyName == nil || *nbResource.Status.PolicyName != nbResource.Spec.PolicyName {
|
||||
nbResource.Status.PolicyName = &nbResource.Spec.PolicyName
|
||||
}
|
||||
if nbResource.Status.PolicyName == nil || *nbResource.Status.PolicyName != nbResource.Spec.PolicyName {
|
||||
nbResource.Status.PolicyName = &nbResource.Spec.PolicyName
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -510,5 +635,18 @@ func (r *NBResourceReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||||
For(&netbirdiov1.NBResource{}).
|
||||
Named("nbresource").
|
||||
Watches(&netbirdiov1.NBGroup{}, handler.EnqueueRequestForOwner(r.Scheme, mgr.GetRESTMapper(), &netbirdiov1.NBResource{})).
|
||||
Watches(&netbirdiov1.NBPolicy{}, handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, obj client.Object) []reconcile.Request {
|
||||
if v, ok := obj.GetAnnotations()["netbird.io/generated-by"]; ok {
|
||||
return []reconcile.Request{
|
||||
{
|
||||
NamespacedName: types.NamespacedName{
|
||||
Namespace: strings.Split(v, "/")[0],
|
||||
Name: strings.Split(v, "/")[1],
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})).
|
||||
Complete(r)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user