Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
85 changes: 58 additions & 27 deletions internal/kube/site/extended_bindings.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,17 +45,18 @@ func NewExtendedBindings(controller *watchers.EventProcessor, profilePath string
return eb
}

func (a *ExtendedBindings) init(context BindingContext, config *qdr.RouterConfig) {
func (a *ExtendedBindings) init(context BindingContext, config *qdr.RouterConfig) error {
a.context = context
if a.mapping == nil {
a.mapping = qdr.RecoverPortMapping(config)
}
a.exposed = ExposedPorts{}
a.selectors = map[string]TargetSelection{}
a.bindings.SetBindingEventHandler(a)
err := a.bindings.SetBindingEventHandler(a)
a.bindings.SetConnectorConfiguration(a.updateBridgeConfigForConnector)
a.bindings.SetListenerConfiguration(a.updateBridgeConfigForListener)
a.bindings.SetMultiKeyListenerConfiguration(a.updateBridgeConfigForMultiKeyListener)
return err
}

func (a *ExtendedBindings) cleanup() {
Expand Down Expand Up @@ -101,15 +102,22 @@ func (a *ExtendedBindings) ConnectorDeleted(connector *skupperv2alpha1.Connector
}
}

func (a *ExtendedBindings) ListenerUpdated(listener *skupperv2alpha1.Listener) {
func (a *ExtendedBindings) ListenerUpdated(listener *skupperv2alpha1.Listener) error {
if err := validateExposedHost(listener.Spec.Host); err != nil {
bindings_logger.Error("Cannot expose listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name),
slog.Any("error", err))
return err
}
allocatedRouterPort, err := a.mapping.GetPortForKey(listener.Name)
if err != nil {
bindings_logger.Error("Unable to get port for listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name),
slog.Any("error", err),
)
return
return err
}
port := Port{
Name: listener.Name,
Expand All @@ -123,12 +131,13 @@ func (a *ExtendedBindings) ListenerUpdated(listener *skupperv2alpha1.Listener) {
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name),
slog.Any("error", err))
} else {
bindings_logger.Info("Exposed listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name))
return err
}
bindings_logger.Info("Exposed listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name))
}
return nil
}

func (a *ExtendedBindings) GetExposedPortSet(host string) *ExposedPortSet {
Expand All @@ -138,27 +147,31 @@ func (a *ExtendedBindings) GetExposedPortSet(host string) *ExposedPortSet {
return nil
}

func (a *ExtendedBindings) ListenerDeleted(listener *skupperv2alpha1.Listener) {
func (a *ExtendedBindings) ListenerDeleted(listener *skupperv2alpha1.Listener) error {
if exposed := a.exposed.Unexpose(listener.Spec.Host, listener.Name); exposed != nil {
a.mapping.ReleasePortForKey(listener.Name)
if exposed.empty() {
if err := a.context.Unexpose(listener.Spec.Host); err != nil {
//TODO: write error to listener status
bindings_logger.Error("Error unexposing service after deleting listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name),
slog.Any("error", err))
return err
}
} else {
if err := a.context.Expose(exposed); err != nil {
//TODO: write error to listener status
bindings_logger.Error("Error re-exposing service after deleting listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name),
slog.Any("error", err))
} else {
bindings_logger.Info("Re-exposed service after deleting listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name))
return err
}
bindings_logger.Info("Re-exposed service after deleting listener",
slog.String("namespace", listener.Namespace),
slog.String("name", listener.Name))
}
}
return nil
}

func (a *ExtendedBindings) updateBridgeConfigForConnector(siteId string, connector *skupperv2alpha1.Connector, config *qdr.BridgeConfig) {
Expand Down Expand Up @@ -213,15 +226,22 @@ func multiKeyListenerPortName(name string) string {
return "multiaddress-" + name
}

func (a *ExtendedBindings) multiKeyListenerUpdated(mkl *skupperv2alpha1.MultiKeyListener) {
func (a *ExtendedBindings) multiKeyListenerUpdated(mkl *skupperv2alpha1.MultiKeyListener) error {
if err := validateExposedHost(mkl.Spec.Host); err != nil {
bindings_logger.Error("Cannot expose multikeylistener",
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name),
slog.Any("error", err))
return err
}
allocatedRouterPort, err := a.mapping.GetPortForKey(multiKeyListenerPortName(mkl.Name))
if err != nil {
bindings_logger.Error("Unable to get port for multikeylistener",
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name),
slog.Any("error", err),
)
return
return err
}
port := Port{
Name: multiKeyListenerPortName(mkl.Name),
Expand All @@ -235,19 +255,20 @@ func (a *ExtendedBindings) multiKeyListenerUpdated(mkl *skupperv2alpha1.MultiKey
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name),
slog.Any("error", err))
return
return err
}
bindings_logger.Info("Exposed multikeylistener",
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name))
}
return nil
}

func (a *ExtendedBindings) multiKeyListenerDeleted(mkl *skupperv2alpha1.MultiKeyListener) {
func (a *ExtendedBindings) multiKeyListenerDeleted(mkl *skupperv2alpha1.MultiKeyListener) error {

exposed := a.exposed.Unexpose(mkl.Spec.Host, multiKeyListenerPortName(mkl.Name))
if exposed == nil {
return
return nil
}
a.mapping.ReleasePortForKey(multiKeyListenerPortName(mkl.Name))
if exposed.empty() {
Expand All @@ -256,19 +277,21 @@ func (a *ExtendedBindings) multiKeyListenerDeleted(mkl *skupperv2alpha1.MultiKey
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name),
slog.Any("error", err))
return err
}
return
return nil
}
if err := a.context.Expose(exposed); err != nil {
bindings_logger.Error("Error re-exposing service after deleting multikeylistener",
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name),
slog.Any("error", err))
return
return err
}
bindings_logger.Info("Re-exposed service after deleting multikeylistener",
slog.String("namespace", mkl.Namespace),
slog.String("name", mkl.Name))
return nil
}

func (b *ExtendedBindings) SetListenerConfiguration(configuration site.ListenerConfiguration) {
Expand All @@ -279,8 +302,8 @@ func (b *ExtendedBindings) SetConnectorConfiguration(configuration site.Connecto
b.bindings.SetConnectorConfiguration(configuration)
}

func (b *ExtendedBindings) SetBindingEventHandler(handler site.BindingEventHandler) {
b.bindings.SetBindingEventHandler(handler)
func (b *ExtendedBindings) SetBindingEventHandler(handler site.BindingEventHandler) error {
return b.bindings.SetBindingEventHandler(handler)
}

func (b *ExtendedBindings) UpdateConnector(name string, connector *skupperv2alpha1.Connector) qdr.ConfigUpdate {
Expand Down Expand Up @@ -321,7 +344,11 @@ func (b *ExtendedBindings) UpdateListener(name string, listener *skupperv2alpha1
}
b.listenerHosts[name] = listener.Spec.Host
}
if b.bindings.UpdateListener(name, listener) != nil {
update, err := b.bindings.UpdateListener(name, listener)
if err != nil {
errs = append(errs, err)
}
if update != nil {
updateConfig = true
}
if !updateConfig {
Expand All @@ -344,13 +371,17 @@ func (b *ExtendedBindings) UpdateMultiKeyListener(name string, mkl *skupperv2alp
}
}
b.multiKeyListenerHosts[name] = mkl.Spec.Host
b.multiKeyListenerUpdated(mkl)
if err := b.multiKeyListenerUpdated(mkl); err != nil {
errs = append(errs, err)
}
} else {
// Deletion case
if previousHost, ok := b.multiKeyListenerHosts[name]; ok {
existingMkl := b.bindings.GetMultiKeyListener(name)
if existingMkl != nil {
b.multiKeyListenerDeleted(existingMkl)
if err := b.multiKeyListenerDeleted(existingMkl); err != nil {
errs = append(errs, err)
}
}
delete(b.multiKeyListenerHosts, name)
if exposed := b.exposed.Unexpose(previousHost, multiKeyListenerPortName(name)); exposed != nil && exposed.empty() {
Expand Down
13 changes: 13 additions & 0 deletions internal/kube/site/ports.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,23 @@
package site

import (
"fmt"
"strings"

corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/apimachinery/pkg/util/validation"
)

// validateExposedHost checks that a host can be used as the name of the
// Service through which it is exposed in the namespace.
func validateExposedHost(host string) error {
if errs := validation.IsDNS1035Label(host); len(errs) > 0 {
return fmt.Errorf("Invalid host %q: %s", host, strings.Join(errs, "; "))
}
return nil
}

type Port struct {
Name string
Port int
Expand Down
20 changes: 15 additions & 5 deletions internal/kube/site/site.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,7 +226,12 @@ func (s *Site) reconcile(siteDef *skupperv2alpha1.Site, inRecovery bool) error {
}
s.initialised = true
s.currentGroups = s.groups()
s.bindings.init(s, routerConfig)
if err := s.bindings.init(s, routerConfig); err != nil {
s.logger.Error("Error exposing existing listeners",
slog.String("namespace", siteDef.Namespace),
slog.String("name", siteDef.Name),
slog.Any("error", err))
}
s.setBindingsConfiguredStatus(nil)
s.checkSecuredAccess()
} else if len(s.currentGroups) != len(s.groups()) {
Expand Down Expand Up @@ -1197,8 +1202,8 @@ func (s *Site) CheckListener(name string, listener *skupperv2alpha1.Listener, sv
return stderrors.Join(err1, s.updateRouterConfig(update))
}
if update == nil {
if tlsErr != nil {
return s.updateListenerStatus(listener, tlsErr)
if err := stderrors.Join(tlsErr, err1); err != nil {
return s.updateListenerStatus(listener, err)
}
return nil
}
Expand All @@ -1215,7 +1220,10 @@ func (s *Site) CheckMultiKeyListener(name string, mkl *skupperv2alpha1.MultiKeyL
}
update, err1 := s.bindings.UpdateMultiKeyListener(name, mkl)
if update == nil {
return nil
if mkl == nil || err1 == nil {
return err1
}
return s.updateMultiKeyListenerStatus(mkl, err1)
}
err2 := s.updateRouterConfig(update)
if mkl == nil {
Expand All @@ -1236,7 +1244,9 @@ func (s *Site) updateMultiKeyListenerStatus(mkl *skupperv2alpha1.MultiKeyListene

func (s *Site) setBindingsConfiguredStatus(err error) {
lf := func(listener *skupperv2alpha1.Listener) *skupperv2alpha1.Listener {
if listener.SetConfigured(nil) {
// a listener whose host cannot be exposed as a service is not configured,
// regardless of the state of the site as a whole
if listener.SetConfigured(validateExposedHost(listener.Spec.Host)) {
updated, err := s.clients.GetSkupperClient().SkupperV2alpha1().Listeners(listener.ObjectMeta.Namespace).UpdateStatus(context.TODO(), listener, metav1.UpdateOptions{})
if err == nil {
return updated
Expand Down
72 changes: 70 additions & 2 deletions internal/kube/site/site_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"context"
"log/slog"
"maps"
"strings"
"testing"

"github.com/skupperproject/skupper/internal/kube/certificates"
Expand All @@ -19,6 +20,7 @@ import (
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
Expand Down Expand Up @@ -403,13 +405,13 @@ func TestSite_CheckListener(t *testing.T) {
RoutingKey: "backend",
Port: 8080,
Type: "tcp",
Host: "1.2.3.4",
Host: "backend",
},
},
svcExists: false,
leadListenersIn: map[string]string{},
leadListenersOut: map[string]string{
"listener1": "1.2.3.4/8080",
"listener1": "backend/8080",
},
},
skupperObjects: []runtime.Object{
Expand Down Expand Up @@ -773,6 +775,72 @@ func TestSite_CheckListener(t *testing.T) {
}
}

// a host that cannot be used as the name of the service exposing it must be
// reported through the status of the listener, not just logged
func TestSite_CheckListenerWithInvalidHost(t *testing.T) {
listener := &skupperv2alpha1.Listener{
ObjectMeta: metav1.ObjectMeta{
Name: "listener1",
Namespace: "test",
},
Spec: skupperv2alpha1.ListenerSpec{
RoutingKey: "backend",
Port: 8080,
Type: "tcp",
Host: "{name}",
},
}
s, err := newSiteMocks("test", nil, []runtime.Object{listener.DeepCopy()}, "", false)
assert.Assert(t, err)
s.initialised = true
assert.Assert(t, createRouterConfigMock(s))

assert.Assert(t, s.CheckListener(listener.Name, listener, false))

condition := meta.FindStatusCondition(listener.Status.Conditions, skupperv2alpha1.CONDITION_TYPE_CONFIGURED)
assert.Assert(t, condition != nil)
assert.Equal(t, condition.Status, metav1.ConditionFalse)
assert.Assert(t, strings.Contains(listener.Status.Message, "Invalid host"), listener.Status.Message)
assert.Equal(t, listener.Status.StatusType, skupperv2alpha1.StatusError)

// no service should have been created for the invalid host
_, err = s.clients.GetKubeClient().CoreV1().Services("test").Get(context.Background(), listener.Spec.Host, metav1.GetOptions{})
assert.Assert(t, errors.IsNotFound(err))
}

func TestSite_CheckMultiKeyListenerWithInvalidHost(t *testing.T) {
mkl := &skupperv2alpha1.MultiKeyListener{
ObjectMeta: metav1.ObjectMeta{
Name: "backend-tcp",
Namespace: "test",
},
Spec: skupperv2alpha1.MultiKeyListenerSpec{
Port: 8080,
Host: "{name}",
Strategy: skupperv2alpha1.MultiKeyListenerStrategy{
Weighted: &skupperv2alpha1.WeightedStrategySpec{
RoutingKeys: map[string]uint{"east-backend": 80},
},
},
},
}
s, err := newSiteMocks("test", nil, []runtime.Object{mkl.DeepCopy()}, "", false)
assert.Assert(t, err)
s.initialised = true
assert.Assert(t, createRouterConfigMock(s))

assert.Assert(t, s.CheckMultiKeyListener(mkl.Name, mkl))

condition := meta.FindStatusCondition(mkl.Status.Conditions, skupperv2alpha1.CONDITION_TYPE_CONFIGURED)
assert.Assert(t, condition != nil)
assert.Equal(t, condition.Status, metav1.ConditionFalse)
assert.Assert(t, strings.Contains(mkl.Status.Message, "Invalid host"), mkl.Status.Message)
assert.Equal(t, mkl.Status.StatusType, skupperv2alpha1.StatusError)

_, err = s.clients.GetKubeClient().CoreV1().Services("test").Get(context.Background(), mkl.Spec.Host, metav1.GetOptions{})
assert.Assert(t, errors.IsNotFound(err))
}

func TestSite_CheckConnector(t *testing.T) {
type args struct {
name string
Expand Down
Loading
Loading