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
12 changes: 12 additions & 0 deletions pkg/apis/serving/v1beta1/domainmapping_lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,18 @@ func (dms *DomainMappingStatus) MarkIngressNotConfigured() {
"IngressNotConfigured", "Ingress has not yet been reconciled.")
}

// MarkTargetIngressNotConfigured marks DomainMappingConditionIngressReady unknown.
func (dms *DomainMappingStatus) MarkTargetIngressNotConfigured(message string) {
domainMappingCondSet.Manage(dms).MarkUnknown(DomainMappingConditionIngressReady,
"IngressNotConfigured", message)
}

// MarkTargetNotOwned marks DomainMappingConditionIngressReady false.
func (dms *DomainMappingStatus) MarkTargetNotOwned(message string) {
domainMappingCondSet.Manage(dms).MarkFalse(DomainMappingConditionIngressReady,
"NotOwned", message)
}

// MarkDomainClaimed updates the DomainMappingConditionDomainClaimed condition
// to indicate that the domain was successfully claimed.
func (dms *DomainMappingStatus) MarkDomainClaimed() {
Expand Down
51 changes: 51 additions & 0 deletions pkg/apis/serving/v1beta1/domainmapping_lifecycle_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,57 @@ func TestReferenceResolvedCondition(t *testing.T) {
apistest.CheckConditionFailed(dms, DomainMappingConditionReady, t)
}

func TestTargetIngressConditions(t *testing.T) {
tests := []struct {
name string
mark func(*DomainMappingStatus)
wantStatus corev1.ConditionStatus
wantReason string
wantMsg string
}{
{
name: "target ingress not configured",
mark: func(dms *DomainMappingStatus) {
dms.MarkTargetIngressNotConfigured("Waiting for target Ingress.")
},
wantStatus: corev1.ConditionUnknown,
wantReason: "IngressNotConfigured",
wantMsg: "Waiting for target Ingress.",
},
{
name: "target resource not owned",
mark: func(dms *DomainMappingStatus) {
dms.MarkTargetNotOwned("Route does not own Ingress.")
},
wantStatus: corev1.ConditionFalse,
wantReason: "NotOwned",
wantMsg: "Route does not own Ingress.",
},
}

for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
dms := &DomainMappingStatus{}
dms.InitializeConditions()
test.mark(dms)

got := dms.GetCondition(DomainMappingConditionIngressReady)
if got == nil {
t.Fatal("IngressReady condition is nil")
}
if got.Status != test.wantStatus {
t.Errorf("Status = %q, want %q", got.Status, test.wantStatus)
}
if got.Reason != test.wantReason {
t.Errorf("Reason = %q, want %q", got.Reason, test.wantReason)
}
if got.Message != test.wantMsg {
t.Errorf("Message = %q, want %q", got.Message, test.wantMsg)
}
})
}
}

func TestCertificateNotReady(t *testing.T) {
dms := &DomainMappingStatus{}

Expand Down
27 changes: 26 additions & 1 deletion pkg/reconciler/domainmapping/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"context"

"k8s.io/client-go/tools/cache"
netv1alpha1 "knative.dev/networking/pkg/apis/networking/v1alpha1"
netclient "knative.dev/networking/pkg/client/injection/client"
certificateinformer "knative.dev/networking/pkg/client/injection/informers/networking/v1alpha1/certificate"
domainclaiminformer "knative.dev/networking/pkg/client/injection/informers/networking/v1alpha1/clusterdomainclaim"
Expand All @@ -29,7 +30,10 @@ import (
"knative.dev/pkg/controller"
"knative.dev/pkg/logging"
"knative.dev/pkg/resolver"
servingv1 "knative.dev/serving/pkg/apis/serving/v1"
"knative.dev/serving/pkg/apis/serving/v1beta1"
routeinformer "knative.dev/serving/pkg/client/injection/informers/serving/v1/route"
serviceinformer "knative.dev/serving/pkg/client/injection/informers/serving/v1/service"
"knative.dev/serving/pkg/client/injection/informers/serving/v1beta1/domainmapping"
kindreconciler "knative.dev/serving/pkg/client/injection/reconciler/serving/v1beta1/domainmapping"
"knative.dev/serving/pkg/reconciler/domainmapping/config"
Expand All @@ -42,11 +46,15 @@ func NewController(ctx context.Context, cmw configmap.Watcher) *controller.Impl
domainmappingInformer := domainmapping.Get(ctx)
ingressInformer := ingressinformer.Get(ctx)
domainClaimInformer := domainclaiminformer.Get(ctx)
routeInformer := routeinformer.Get(ctx)
serviceInformer := serviceinformer.Get(ctx)

r := &Reconciler{
certificateLister: certificateInformer.Lister(),
ingressLister: ingressInformer.Lister(),
domainClaimLister: domainClaimInformer.Lister(),
routeLister: routeInformer.Lister(),
serviceLister: serviceInformer.Lister(),
netclient: netclient.Get(ctx),
}

Expand All @@ -71,7 +79,24 @@ func NewController(ctx context.Context, cmw configmap.Watcher) *controller.Impl
certificateInformer.Informer().AddEventHandler(handleControllerOf)
ingressInformer.Informer().AddEventHandler(handleControllerOf)

r.resolver = resolver.NewURIResolverFromTracker(ctx, impl.Tracker)
r.tracker = impl.Tracker
r.resolver = resolver.NewURIResolverFromTracker(ctx, r.tracker)

// Track Route and Ingress changes because their HTTP paths are copied into
// DomainMapping Ingresses. The resolver already tracks referenced Services.
// Informer events may omit TypeMeta, so populate it before notifying the tracker.
routeInformer.Informer().AddEventHandler(controller.HandleAll(
controller.EnsureTypeMeta(r.tracker.OnChanged, servingv1.SchemeGroupVersion.WithKind("Route")),
))
ingressInformer.Informer().AddEventHandler(cache.FilteringResourceEventHandler{
FilterFunc: controller.FilterController(&servingv1.Route{}),
Handler: controller.HandleAll(
controller.EnsureTypeMeta(r.tracker.OnChanged, netv1alpha1.SchemeGroupVersion.WithKind("Ingress")),
),
})
domainmappingInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
DeleteFunc: r.tracker.OnDeletedObserver,
})

return impl
}
142 changes: 139 additions & 3 deletions pkg/reconciler/domainmapping/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ package domainmapping

import (
"context"
"errors"
"fmt"
"slices"
"sort"
Expand All @@ -31,6 +32,7 @@ import (
"k8s.io/apimachinery/pkg/api/equality"
apierrs "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"

netapi "knative.dev/networking/pkg/apis/networking"
netv1alpha1 "knative.dev/networking/pkg/apis/networking/v1alpha1"
Expand All @@ -44,22 +46,30 @@ import (
"knative.dev/pkg/network"
"knative.dev/pkg/reconciler"
"knative.dev/pkg/resolver"
v1 "knative.dev/serving/pkg/apis/serving/v1"
"knative.dev/pkg/tracker"
"knative.dev/serving/pkg/apis/serving"
servingv1 "knative.dev/serving/pkg/apis/serving/v1"
"knative.dev/serving/pkg/apis/serving/v1beta1"
domainmappingreconciler "knative.dev/serving/pkg/client/injection/reconciler/serving/v1beta1/domainmapping"
servinglisters "knative.dev/serving/pkg/client/listers/serving/v1"
servingnetworking "knative.dev/serving/pkg/networking"
"knative.dev/serving/pkg/reconciler/domainmapping/config"
"knative.dev/serving/pkg/reconciler/domainmapping/resources"
routeresources "knative.dev/serving/pkg/reconciler/route/resources"
routenames "knative.dev/serving/pkg/reconciler/route/resources/names"
servicenames "knative.dev/serving/pkg/reconciler/service/resources/names"
)

// Reconciler implements controller.Reconciler for DomainMapping resources.
type Reconciler struct {
certificateLister networkinglisters.CertificateLister
ingressLister networkinglisters.IngressLister
domainClaimLister networkinglisters.ClusterDomainClaimLister
serviceLister servinglisters.ServiceLister
routeLister servinglisters.RouteLister
netclient netclientset.Interface
resolver *resolver.URIResolver
tracker tracker.Interface
}

// Check that our Reconciler implements Interface
Expand Down Expand Up @@ -116,12 +126,32 @@ func (r *Reconciler) ReconcileKind(ctx context.Context, dm *v1beta1.DomainMappin
return err
}

// For Knative targets, preserve Route traffic instead of sending everything
// to the Service resolved from the Addressable URL.
targetKind, preserveRouteTraffic := servingReferenceKind(dm.Spec.Ref)

// Resolve the spec.Ref to a URI following the Addressable contract.
targetHost, targetBackendSvc, err := r.resolveRef(ctx, dm)
if err != nil {
if !preserveRouteTraffic {
return err
}
dm.Status.MarkIngressNotConfigured()
return err
}

var routePaths []netv1alpha1.HTTPIngressPath
if preserveRouteTraffic {
var found bool
routePaths, found, err = r.routeHTTPPaths(dm, targetKind, targetHost)
if err != nil {
return err
}
if !found {
return nil
}
}

// HTTPOption can be set via annotations or in the config map.
httpOption, err := servingnetworking.GetHTTPOption(ctx, config.FromContext(ctx).Network, dm.GetAnnotations())
if err != nil {
Expand All @@ -130,7 +160,12 @@ func (r *Reconciler) ReconcileKind(ctx context.Context, dm *v1beta1.DomainMappin

// Reconcile the Ingress resource corresponding to the requested Mapping.
logger.Debugf("Mapping %s to ref %s/%s (host: %q, svc: %q)", url, dm.Spec.Ref.Namespace, dm.Spec.Ref.Name, targetHost, targetBackendSvc)
desired := resources.MakeIngress(dm, targetBackendSvc, targetHost, ingressClass, httpOption, tls, acmeChallenges...)
var desired *netv1alpha1.Ingress
if preserveRouteTraffic {
desired = resources.MakeIngressWithHTTPPaths(dm, routePaths, ingressClass, httpOption, tls, acmeChallenges...)
} else {
desired = resources.MakeIngress(dm, targetBackendSvc, targetHost, ingressClass, httpOption, tls, acmeChallenges...)
}
ingress, err := r.reconcileIngress(ctx, dm, desired)
if err != nil {
return err
Expand Down Expand Up @@ -205,7 +240,7 @@ func (r *Reconciler) tls(ctx context.Context, dm *v1beta1.DomainMapping) ([]netv
}

if !externalDomainTLSEnabled(ctx, dm) {
dm.Status.MarkTLSNotEnabled(v1.ExternalDomainTLSNotEnabledMessage)
dm.Status.MarkTLSNotEnabled(servingv1.ExternalDomainTLSNotEnabledMessage)
return nil, nil, nil
}

Expand Down Expand Up @@ -317,6 +352,107 @@ func (r *Reconciler) resolveRef(ctx context.Context, dm *v1beta1.DomainMapping)
return resolved.Host, parts[0], nil
}

// routeHTTPPaths gets the first matching cluster-local paths for a Service or
// Route target. If no expected source can be found, it updates IngressReady and
// returns found=false. The source Ingress is responsible for validating paths.
// The returned paths belong to the informer cache and must not be mutated.
func (r *Reconciler) routeHTTPPaths(dm *v1beta1.DomainMapping, kind, targetHost string) (paths []netv1alpha1.HTTPIngressPath, found bool, err error) {
ref := dm.Spec.Ref

if r.tracker == nil {
return nil, false, errors.New("DomainMapping target tracker is not configured")
}

routeName := ref.Name
var service *servingv1.Service
if kind == "Service" {
service, err = r.serviceLister.Services(ref.Namespace).Get(ref.Name)
if apierrs.IsNotFound(err) {
dm.Status.MarkTargetIngressNotConfigured(fmt.Sprintf("Waiting for target Service %s/%s to be observed.", ref.Namespace, ref.Name))
return nil, false, nil
} else if err != nil {
return nil, false, fmt.Errorf("getting target Service %s/%s: %w", ref.Namespace, ref.Name, err)
}
routeName = servicenames.Route(service)
}

if err := r.tracker.TrackReference(tracker.Reference{
APIVersion: servingv1.SchemeGroupVersion.String(),
Kind: "Route",
Namespace: ref.Namespace,
Name: routeName,
}, dm); err != nil {
return nil, false, fmt.Errorf("tracking target Route: %w", err)
}

route, err := r.routeLister.Routes(ref.Namespace).Get(routeName)
if apierrs.IsNotFound(err) {
dm.Status.MarkTargetIngressNotConfigured(fmt.Sprintf("Waiting for target Route %s/%s to be created.", ref.Namespace, routeName))
return nil, false, nil
} else if err != nil {
return nil, false, fmt.Errorf("getting target Route %s/%s: %w", ref.Namespace, routeName, err)
}

if service != nil && !metav1.IsControlledBy(route, service) {
dm.Status.MarkTargetNotOwned(fmt.Sprintf("Service %s/%s does not own Route %s/%s.", service.Namespace, service.Name, route.Namespace, route.Name))
return nil, false, nil
}

ingressName := routenames.Ingress(route)
if err := r.tracker.TrackReference(tracker.Reference{
APIVersion: netv1alpha1.SchemeGroupVersion.String(),
Kind: "Ingress",
Namespace: route.Namespace,
Name: ingressName,
}, dm); err != nil {
return nil, false, fmt.Errorf("tracking target Ingress: %w", err)
}

ingress, err := r.ingressLister.Ingresses(route.Namespace).Get(ingressName)
if apierrs.IsNotFound(err) {
dm.Status.MarkTargetIngressNotConfigured(fmt.Sprintf("Waiting for target Route %s/%s to configure Ingress %s/%s.", route.Namespace, route.Name, route.Namespace, ingressName))
return nil, false, nil
} else if err != nil {
return nil, false, fmt.Errorf("getting target Ingress %s/%s: %w", route.Namespace, ingressName, err)
}
if !metav1.IsControlledBy(ingress, route) {
dm.Status.MarkTargetNotOwned(fmt.Sprintf("Route %s/%s does not own Ingress %s/%s.", route.Namespace, route.Name, ingress.Namespace, ingress.Name))
return nil, false, nil
}

for i := range ingress.Spec.Rules {
rule := &ingress.Spec.Rules[i]
if rule.Visibility == netv1alpha1.IngressVisibilityClusterLocal && slices.Contains(rule.Hosts, targetHost) {
if rule.HTTP == nil {
return nil, false, fmt.Errorf("target Ingress %s/%s has a matching cluster-local rule without HTTP", ingress.Namespace, ingress.Name)
}
return rule.HTTP.Paths, true, nil
}
}

dm.Status.MarkTargetIngressNotConfigured(fmt.Sprintf("Ingress %s/%s has no cluster-local rule for host %q.", ingress.Namespace, ingress.Name, targetHost))
return nil, false, nil
}

func servingReferenceKind(ref duckv1.KReference) (string, bool) {
group := ref.Group
if ref.APIVersion != "" {
gv, err := schema.ParseGroupVersion(ref.APIVersion)
if err != nil {
return "", false
}
group = gv.Group
}
if group != serving.GroupName {
return "", false
}
switch ref.Kind {
case "Service", "Route":
return ref.Kind, true
}
return "", false
}

func (r *Reconciler) reconcileDomainClaim(ctx context.Context, dm *v1beta1.DomainMapping) error {
dc, err := r.domainClaimLister.Get(dm.Name)
if err != nil && !apierrs.IsNotFound(err) {
Expand Down
Loading
Loading