Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
5 changes: 4 additions & 1 deletion internal/adc/translator/apisixconsumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,10 @@ func (t *Translator) TranslateApisixConsumer(tctx *provider.TranslateContext, ac
if !plugin.Enable {
continue
}
config := t.buildPluginConfig(plugin, ac.Namespace, tctx.Secrets)
config, err := t.buildPluginConfig(plugin, ac.Namespace, tctx.Secrets)
if err != nil {
return nil, err
}
plugins[plugin.Name] = config
}

Expand Down
47 changes: 32 additions & 15 deletions internal/adc/translator/apisixroute.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,10 @@ func (t *Translator) TranslateApisixRoute(tctx *provider.TranslateContext, ar *a

func (t *Translator) translateHTTPRule(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, rule apiv2.ApisixRouteHTTP, ruleIndex int) (*adc.Service, error) {
timeout := t.buildTimeout(rule)
plugins := t.buildPlugins(tctx, ar, rule)
plugins, err := t.buildPlugins(tctx, ar, rule)
if err != nil {
return nil, err
}

vars, err := rule.Match.NginxVars.ToVars()
if err != nil {
Expand Down Expand Up @@ -91,24 +94,28 @@ func (t *Translator) buildTimeout(rule apiv2.ApisixRouteHTTP) *adc.Timeout {
}
}

func (t *Translator) buildPlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, rule apiv2.ApisixRouteHTTP) adc.Plugins {
func (t *Translator) buildPlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, rule apiv2.ApisixRouteHTTP) (adc.Plugins, error) {
plugins := make(adc.Plugins)

// Load plugins from referenced PluginConfig
t.loadPluginConfigPlugins(tctx, ar, rule, plugins)
if err := t.loadPluginConfigPlugins(tctx, ar, rule, plugins); err != nil {
return nil, err
}

// Apply plugins from the route itself
t.loadRoutePlugins(tctx, ar, rule.Plugins, plugins)
if err := t.loadRoutePlugins(tctx, ar, rule.Plugins, plugins); err != nil {
return nil, err
}

// Add authentication plugins
t.addAuthenticationPlugins(rule, plugins)

return plugins
return plugins, nil
}

func (t *Translator) loadPluginConfigPlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, rule apiv2.ApisixRouteHTTP, plugins adc.Plugins) {
func (t *Translator) loadPluginConfigPlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, rule apiv2.ApisixRouteHTTP, plugins adc.Plugins) error {
if rule.PluginConfigName == "" {
return
return nil
}

pcNamespace := ar.Namespace
Expand All @@ -119,33 +126,41 @@ func (t *Translator) loadPluginConfigPlugins(tctx *provider.TranslateContext, ar
pcKey := types.NamespacedName{Namespace: pcNamespace, Name: rule.PluginConfigName}
pc, ok := tctx.ApisixPluginConfigs[pcKey]
if !ok || pc == nil {
return
return nil
}

for _, plugin := range pc.Spec.Plugins {
if !plugin.Enable {
continue
}
config := t.buildPluginConfig(plugin, pc.Namespace, tctx.Secrets)
config, err := t.buildPluginConfig(plugin, pc.Namespace, tctx.Secrets)
if err != nil {
return err
}
plugins[plugin.Name] = config
}
return nil
}

func (t *Translator) loadRoutePlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, routePlugins []apiv2.ApisixRoutePlugin, plugins adc.Plugins) {
func (t *Translator) loadRoutePlugins(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, routePlugins []apiv2.ApisixRoutePlugin, plugins adc.Plugins) error {
for _, plugin := range routePlugins {
if !plugin.Enable {
continue
}
config := t.buildPluginConfig(plugin, ar.Namespace, tctx.Secrets)
config, err := t.buildPluginConfig(plugin, ar.Namespace, tctx.Secrets)
if err != nil {
return err
}
plugins[plugin.Name] = config
}
return nil
}

func (t *Translator) buildPluginConfig(plugin apiv2.ApisixRoutePlugin, namespace string, secrets map[types.NamespacedName]*corev1.Secret) map[string]any {
func (t *Translator) buildPluginConfig(plugin apiv2.ApisixRoutePlugin, namespace string, secrets map[types.NamespacedName]*corev1.Secret) (map[string]any, error) {
config := make(map[string]any)
if len(plugin.Config.Raw) > 0 {
if err := json.Unmarshal(plugin.Config.Raw, &config); err != nil {
t.Log.Error(err, "failed to unmarshal plugin config")
return nil, fmt.Errorf("failed to unmarshal config of plugin %s: %w", plugin.Name, err)
}
}
if plugin.SecretRef != "" {
Expand All @@ -155,7 +170,7 @@ func (t *Translator) buildPluginConfig(plugin apiv2.ApisixRoutePlugin, namespace
}
}
}
return config
return config, nil
}

func (t *Translator) addAuthenticationPlugins(rule apiv2.ApisixRouteHTTP, plugins adc.Plugins) {
Expand Down Expand Up @@ -473,7 +488,9 @@ func (t *Translator) translateApisixRouteBackendResolveGranularityEndpoint(tctx
func (t *Translator) translateStreamRule(tctx *provider.TranslateContext, ar *apiv2.ApisixRoute, part apiv2.ApisixRouteStream) (*adc.Service, error) {
// add stream route plugins
plugins := make(adc.Plugins)
t.loadRoutePlugins(tctx, ar, part.Plugins, plugins)
if err := t.loadRoutePlugins(tctx, ar, part.Plugins, plugins); err != nil {
return nil, err
}

sr := adc.NewDefaultStreamRoute()
sr.Name = adc.ComposeStreamRouteName(ar.Namespace, ar.Name, part.Name, part.Protocol)
Expand Down
5 changes: 4 additions & 1 deletion internal/adc/translator/globalrule.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,10 @@ func (t *Translator) TranslateApisixGlobalRule(tctx *provider.TranslateContext,
continue
}

pluginConfig := t.buildPluginConfig(plugin, obj.Namespace, tctx.Secrets)
pluginConfig, err := t.buildPluginConfig(plugin, obj.Namespace, tctx.Secrets)
if err != nil {
return nil, err
}
plugins[plugin.Name] = pluginConfig
}

Expand Down
11 changes: 8 additions & 3 deletions internal/adc/translator/grpcroute.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ func (t *Translator) fillPluginsFromGRPCRouteFilters(
namespace string,
filters []gatewayv1.GRPCRouteFilter,
tctx *provider.TranslateContext,
) {
) error {
for _, filter := range filters {
switch filter.Type {
case gatewayv1.GRPCRouteFilterRequestHeaderModifier:
Expand All @@ -48,9 +48,12 @@ func (t *Translator) fillPluginsFromGRPCRouteFilters(
case gatewayv1.GRPCRouteFilterResponseHeaderModifier:
t.fillPluginFromHTTPResponseHeaderFilter(plugins, filter.ResponseHeaderModifier)
case gatewayv1.GRPCRouteFilterExtensionRef:
t.fillPluginFromExtensionRef(plugins, namespace, filter.ExtensionRef, tctx)
if err := t.fillPluginFromExtensionRef(plugins, namespace, filter.ExtensionRef, tctx); err != nil {
return err
}
}
}
return nil
}

func calculateGRPCRoutePriority(match *gatewayv1.GRPCRouteMatch, ruleIndex int, hosts []string) uint64 {
Expand Down Expand Up @@ -283,7 +286,9 @@ func (t *Translator) TranslateGRPCRoute(tctx *provider.TranslateContext, grpcRou
}
}

t.fillPluginsFromGRPCRouteFilters(service.Plugins, grpcRoute.GetNamespace(), rule.Filters, tctx)
if err := t.fillPluginsFromGRPCRouteFilters(service.Plugins, grpcRoute.GetNamespace(), rule.Filters, tctx); err != nil {
return nil, err
}

matches := rule.Matches
if len(matches) == 0 {
Expand Down
21 changes: 13 additions & 8 deletions internal/adc/translator/httproute.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ func (t *Translator) fillPluginsFromHTTPRouteFilters(
filters []gatewayv1.HTTPRouteFilter,
matches []gatewayv1.HTTPRouteMatch,
tctx *provider.TranslateContext,
) {
) error {
for _, filter := range filters {
switch filter.Type {
case gatewayv1.HTTPRouteFilterRequestHeaderModifier:
Expand All @@ -59,38 +59,41 @@ func (t *Translator) fillPluginsFromHTTPRouteFilters(
case gatewayv1.HTTPRouteFilterResponseHeaderModifier:
t.fillPluginFromHTTPResponseHeaderFilter(plugins, filter.ResponseHeaderModifier)
case gatewayv1.HTTPRouteFilterExtensionRef:
t.fillPluginFromExtensionRef(plugins, namespace, filter.ExtensionRef, tctx)
if err := t.fillPluginFromExtensionRef(plugins, namespace, filter.ExtensionRef, tctx); err != nil {
return err
}
case gatewayv1.HTTPRouteFilterCORS:
t.fillPluginFromHTTPCORSFilter(plugins, filter.CORS)
}
}
return nil
}

func (t *Translator) fillPluginFromExtensionRef(plugins adctypes.Plugins, namespace string, extensionRef *gatewayv1.LocalObjectReference, tctx *provider.TranslateContext) {
func (t *Translator) fillPluginFromExtensionRef(plugins adctypes.Plugins, namespace string, extensionRef *gatewayv1.LocalObjectReference, tctx *provider.TranslateContext) error {
if extensionRef == nil {
return
return nil
}
if extensionRef.Kind == internaltypes.KindPluginConfig {
pluginconfig := tctx.PluginConfigs[types.NamespacedName{
Namespace: namespace,
Name: string(extensionRef.Name),
}]
if pluginconfig == nil {
return
return nil
}
for _, plugin := range pluginconfig.Spec.Plugins {
pluginName := plugin.Name
pluginconfig := make(map[string]any)
if len(plugin.Config.Raw) > 0 {
if err := json.Unmarshal(plugin.Config.Raw, &pluginconfig); err != nil {
t.Log.Error(err, "plugin config unmarshal failed", "plugin", plugin.Name)
continue
return fmt.Errorf("failed to unmarshal config of plugin %s: %w", plugin.Name, err)
}
}
plugins[pluginName] = pluginconfig
}
t.Log.V(1).Info("fill plugin from extension ref", "plugins", plugins)
}
return nil
}

func (t *Translator) fillPluginFromURLRewriteFilter(plugins adctypes.Plugins, urlRewrite *gatewayv1.HTTPURLRewriteFilter, matches []gatewayv1.HTTPRouteMatch) {
Expand Down Expand Up @@ -668,7 +671,9 @@ func (t *Translator) TranslateHTTPRoute(tctx *provider.TranslateContext, httpRou
}
}

t.fillPluginsFromHTTPRouteFilters(service.Plugins, httpRoute.GetNamespace(), rule.Filters, rule.Matches, tctx)
if err := t.fillPluginsFromHTTPRouteFilters(service.Plugins, httpRoute.GetNamespace(), rule.Filters, rule.Matches, tctx); err != nil {
return nil, err
}

matches := rule.Matches
if len(matches) == 0 {
Expand Down
38 changes: 26 additions & 12 deletions internal/adc/translator/ingress.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,11 @@ func (t *Translator) TranslateIngress(

for j, path := range rule.HTTP.Paths {
index := fmt.Sprintf("%d-%d", i, j)
if svc := t.buildServiceFromIngressPath(tctx, obj, config, &path, index, hosts, labels); svc != nil {
svc, err := t.buildServiceFromIngressPath(tctx, obj, config, &path, index, hosts, labels)
if err != nil {
return nil, err
}
if svc != nil {
result.Services = append(result.Services, svc)
}
}
Expand Down Expand Up @@ -147,9 +151,9 @@ func (t *Translator) buildServiceFromIngressPath(
index string,
hosts []string,
labels map[string]string,
) *adctypes.Service {
) (*adctypes.Service, error) {
if path.Backend.Service == nil {
return nil
return nil, nil
}

service := adctypes.NewDefaultService()
Expand All @@ -162,7 +166,10 @@ func (t *Translator) buildServiceFromIngressPath(
protocol := t.resolveIngressUpstream(tctx, obj, config, path.Backend.Service, upstream)
service.Upstream = upstream

route := t.buildRouteFromIngressPath(tctx, obj, path, config, index, labels)
route, err := t.buildRouteFromIngressPath(tctx, obj, path, config, index, labels)
if err != nil {
return nil, err
}
// Check if websocket is enabled via annotation first, then fall back to appProtocol detection
if config != nil && config.EnableWebsocket {
route.EnableWebsocket = ptr.To(true)
Expand All @@ -172,7 +179,7 @@ func (t *Translator) buildServiceFromIngressPath(
service.Routes = []*adctypes.Route{route}

t.fillHTTPRoutePoliciesForIngress(tctx, service.Routes)
return service
return service, nil
}

func (t *Translator) resolveIngressUpstream(
Expand Down Expand Up @@ -260,7 +267,7 @@ func (t *Translator) buildRouteFromIngressPath(
config *IngressConfig,
index string,
labels map[string]string,
) *adctypes.Route {
) (*adctypes.Route, error) {
route := adctypes.NewDefaultRoute()
route.Name = adctypes.ComposeRouteName(obj.Namespace, obj.Name, index)
route.ID = id.GenID(route.Name)
Expand Down Expand Up @@ -306,7 +313,11 @@ func (t *Translator) buildRouteFromIngressPath(
if config != nil {
// check if PluginConfig is specified
if config.PluginConfigName != "" {
route.Plugins = t.loadPluginConfigPluginsForIngress(tctx, obj.Namespace, config.PluginConfigName)
plugins, err := t.loadPluginConfigPluginsForIngress(tctx, obj.Namespace, config.PluginConfigName)
if err != nil {
return nil, err
}
route.Plugins = plugins
}

// apply plugins from annotations
Expand All @@ -321,10 +332,10 @@ func (t *Translator) buildRouteFromIngressPath(
}

route.Uris = uris
return route
return route, nil
}

func (t *Translator) loadPluginConfigPluginsForIngress(tctx *provider.TranslateContext, namespace, pluginConfigName string) adctypes.Plugins {
func (t *Translator) loadPluginConfigPluginsForIngress(tctx *provider.TranslateContext, namespace, pluginConfigName string) (adctypes.Plugins, error) {
plugins := make(adctypes.Plugins)

pcKey := types.NamespacedName{
Expand All @@ -333,18 +344,21 @@ func (t *Translator) loadPluginConfigPluginsForIngress(tctx *provider.TranslateC
}
pc, ok := tctx.ApisixPluginConfigs[pcKey]
if !ok || pc == nil {
return plugins
return plugins, nil
}

for _, plugin := range pc.Spec.Plugins {
if !plugin.Enable {
continue
}
config := t.buildPluginConfig(plugin, namespace, tctx.Secrets)
config, err := t.buildPluginConfig(plugin, namespace, tctx.Secrets)
if err != nil {
return nil, err
}
plugins[plugin.Name] = config
}

return plugins
return plugins, nil
}

// translateEndpointSliceForIngress create upstream nodes from EndpointSlice
Expand Down
Loading
Loading