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
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