This is an automated email from the ASF dual-hosted git repository.
Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git
The following commit(s) were added to refs/heads/master by this push:
new 6fbffaf43f test: add unit tests for k8s ingress parsers and
IngressReconciler (#6679) (#7175)
6fbffaf43f is described below
commit 6fbffaf43f3b25b3ae22f02ab39da226a5ae1d3e
Author: wy471x <[email protected]>
AuthorDate: Thu Oct 1 20:46:38 2026 +0800
test: add unit tests for k8s ingress parsers and IngressReconciler (#6679)
(#7175)
* test: add unit tests for k8s ingress parsers and IngressReconciler (#6679)
* test: set context path annotation in IngressParser prefix dispatch case
The prefix-path dispatch test passed no context path annotation, so
ContextPathParser skipped the path and produced a config without a
selector, making the CONTEXT_PATH assertion fail.
* test: use ServiceIngressRelation in IngressReconcilerTest
ServiceIngressCache.getIngressName now returns List<ServiceIngressRelation>
instead of List<Pair<String, String>>, which broke test compilation.
---------
Co-authored-by: aias00 <[email protected]>
---
.../apache/shenyu/k8s/IngressReconcilerTest.java | 261 +++++++++++++++++
.../shenyu/k8s/parser/ContextPathParserTest.java | 202 ++++++++++++++
.../apache/shenyu/k8s/parser/GrpcParserTest.java | 275 ++++++++++++++++++
.../shenyu/k8s/parser/IngressParserTest.java | 195 +++++++++++++
.../apache/shenyu/k8s/parser/SofaParserTest.java | 233 ++++++++++++++++
.../shenyu/k8s/parser/WebSocketParserTest.java | 310 +++++++++++++++++++++
6 files changed, 1476 insertions(+)
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/IngressReconcilerTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/IngressReconcilerTest.java
new file mode 100644
index 0000000000..15bd299f13
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/IngressReconcilerTest.java
@@ -0,0 +1,261 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s;
+
+import io.kubernetes.client.extended.controller.reconciler.Request;
+import io.kubernetes.client.extended.controller.reconciler.Result;
+import io.kubernetes.client.informer.SharedIndexInformer;
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.openapi.ApiClient;
+import io.kubernetes.client.openapi.models.CoreV1EndpointPort;
+import io.kubernetes.client.openapi.models.V1EndpointAddress;
+import io.kubernetes.client.openapi.models.V1EndpointSubsetBuilder;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1EndpointsBuilder;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1Secret;
+import io.kubernetes.client.openapi.models.V1Service;
+import org.apache.shenyu.common.config.ssl.ShenyuSniAsyncMapping;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.k8s.cache.IngressCache;
+import org.apache.shenyu.k8s.cache.IngressSelectorCache;
+import org.apache.shenyu.k8s.cache.ServiceIngressCache;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ServiceIngressRelation;
+import org.apache.shenyu.k8s.parser.IngressParser;
+import org.apache.shenyu.k8s.reconciler.IngressReconciler;
+import org.apache.shenyu.k8s.repository.ShenyuCacheRepository;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test for {@link IngressReconciler}.
+ */
+public final class IngressReconcilerTest {
+
+ private static final String NAMESPACE = "reconciler-namespace";
+
+ private static final String INGRESS_NAME = "reconciler-ingress";
+
+ private static final String SERVICE_NAME = "reconciler-service";
+
+ private static final int BACKEND_PORT = 8080;
+
+ private static final int ENDPOINT_PORT = 9090;
+
+ private Indexer<V1Ingress> ingressIndexer;
+
+ private ShenyuCacheRepository shenyuCacheRepository;
+
+ private IngressReconciler ingressReconciler;
+
+ private final List<SelectorData> savedSelectorData = new ArrayList<>();
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void init() {
+ SharedIndexInformer<V1Ingress> ingressInformer =
mock(SharedIndexInformer.class);
+ ingressIndexer = mock(Indexer.class);
+ when(ingressInformer.getIndexer()).thenReturn(ingressIndexer);
+
+ SharedIndexInformer<V1Service> serviceInformer =
mock(SharedIndexInformer.class);
+ when(serviceInformer.getIndexer()).thenReturn(mock(Indexer.class));
+
+ SharedIndexInformer<V1Endpoints> endpointsInformer =
mock(SharedIndexInformer.class);
+ Indexer<V1Endpoints> endpointsIndexer = mock(Indexer.class);
+ when(endpointsInformer.getIndexer()).thenReturn(endpointsIndexer);
+ V1Endpoints endpoints = new V1EndpointsBuilder()
+ .withKind("Endpoints")
+
.withNewMetadata().withNamespace(NAMESPACE).withName(SERVICE_NAME).endMetadata()
+ .withSubsets(new V1EndpointSubsetBuilder()
+ .withAddresses(new V1EndpointAddress().ip("10.0.0.1"))
+ .withPorts(new
CoreV1EndpointPort().port(ENDPOINT_PORT).protocol("TCP"))
+ .build())
+ .build();
+ when(endpointsIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(endpoints);
+
+ shenyuCacheRepository = mock(ShenyuCacheRepository.class);
+
when(shenyuCacheRepository.findRuleDataList(anyString())).thenReturn(Collections.emptyList());
+ doAnswer(invocation -> {
+ savedSelectorData.add(invocation.getArgument(0));
+ return null;
+ }).when(shenyuCacheRepository).saveOrUpdateSelectorData(any());
+
when(shenyuCacheRepository.findSelectorDataList(anyString())).thenAnswer(invocation
-> new ArrayList<>(savedSelectorData));
+
+ SharedIndexInformer<V1Secret> secretInformer =
mock(SharedIndexInformer.class);
+ ingressReconciler = new IngressReconciler(ingressInformer,
secretInformer, shenyuCacheRepository,
+ new ShenyuSniAsyncMapping(), new
IngressParser(serviceInformer, endpointsInformer), mock(ApiClient.class));
+
+ IngressCache.getInstance().remove(NAMESPACE, INGRESS_NAME);
+ IngressSelectorCache.getInstance().remove(NAMESPACE, INGRESS_NAME,
PluginEnum.DIVIDE.getName());
+ IngressSelectorCache.getInstance().remove(NAMESPACE, INGRESS_NAME,
PluginEnum.WEB_SOCKET.getName());
+ ServiceIngressCache.getInstance().removeAllIngressName(NAMESPACE,
SERVICE_NAME);
+ }
+
+ @Test
+ public void
testReconcileNewIngressSavesConfigAndRefreshesUpstreamFromEndpoints() {
+ mockIngress(shenyuAnnotations(new HashMap<>()), "/test", "Exact",
null);
+
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ assertNotNull(IngressCache.getInstance().get(NAMESPACE, INGRESS_NAME));
+ List<ServiceIngressRelation> ingressRelations =
ServiceIngressCache.getInstance().getIngressName(NAMESPACE, SERVICE_NAME);
+ assertNotNull(ingressRelations);
+ assertTrue(ingressRelations.stream().anyMatch(relation ->
relation.isSameIngress(NAMESPACE, INGRESS_NAME)));
+ verify(shenyuCacheRepository).saveOrUpdateRuleData(any());
+
+ // the selector is first saved with the ingress backend port and then
refreshed with the endpoints port
+ assertEquals(2, savedSelectorData.size());
+ assertEquals(PluginEnum.DIVIDE.getName(),
savedSelectorData.get(0).getPluginName());
+ assertThat(savedSelectorData.get(0).getHandle(),
containsString("10.0.0.1:" + BACKEND_PORT));
+ assertThat(savedSelectorData.get(1).getHandle(),
containsString("10.0.0.1:" + ENDPOINT_PORT));
+ assertEquals(savedSelectorData.get(0).getId(),
savedSelectorData.get(1).getId());
+ }
+
+ @Test
+ public void
testReconcileWebSocketIngressEnablesPluginAndUsesWebSocketUpstream() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_WEB_SOCKET_ENABLED, "true");
+ mockIngress(shenyuAnnotations(annotations), "/test", "Prefix", null);
+
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ ArgumentCaptor<PluginData> pluginCaptor =
ArgumentCaptor.forClass(PluginData.class);
+ verify(shenyuCacheRepository,
atLeastOnce()).saveOrUpdatePluginData(pluginCaptor.capture());
+ assertTrue(pluginCaptor.getAllValues().stream()
+ .anyMatch(pluginData ->
PluginEnum.WEB_SOCKET.getName().equals(pluginData.getName())));
+
+ SelectorData savedSelector =
savedSelectorData.get(savedSelectorData.size() - 1);
+ assertEquals(PluginEnum.WEB_SOCKET.getName(),
savedSelector.getPluginName());
+ assertThat(savedSelector.getHandle(),
containsString("\"protocol\":\"ws://\""));
+ assertThat(savedSelector.getHandle(), containsString("10.0.0.1:" +
ENDPOINT_PORT));
+ }
+
+ @Test
+ public void testReconcileSkipsIngressOfAnotherIngressClass() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.K8S_INGRESS_CLASS_ANNOTATION_KEY,
"nginx");
+ mockIngress(annotations, "/test", "Prefix", null);
+
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ verify(shenyuCacheRepository, never()).saveOrUpdateSelectorData(any());
+ assertNull(IngressCache.getInstance().get(NAMESPACE, INGRESS_NAME));
+ }
+
+ @Test
+ public void testReconcileAcceptsIngressClassNameFromSpec() {
+ mockIngress(new HashMap<>(), "/test", "Prefix",
IngressConstants.SHENYU_INGRESS_CLASS);
+
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ verify(shenyuCacheRepository,
atLeastOnce()).saveOrUpdateSelectorData(any());
+ assertNotNull(IngressCache.getInstance().get(NAMESPACE, INGRESS_NAME));
+ }
+
+ @Test
+ public void testReconcileUpdatesConfigWhenIngressChanged() {
+ mockIngress(shenyuAnnotations(new HashMap<>()), "/test", "Exact",
null);
+ ingressReconciler.reconcile(new Request(NAMESPACE, INGRESS_NAME));
+ assertEquals(2, savedSelectorData.size());
+
+ mockIngress(shenyuAnnotations(new HashMap<>()), "/test-changed",
"Exact", null);
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ verify(shenyuCacheRepository, atLeastOnce())
+ .deleteSelectorData(eq(PluginEnum.DIVIDE.getName()),
eq(savedSelectorData.get(0).getId()));
+ // the stale config is deleted and the changed ingress is saved again
with a refreshed upstream handle
+ assertEquals(4, savedSelectorData.size());
+ }
+
+ @Test
+ public void testReconcileDeletedIngressRemovesConfig() {
+ mockIngress(shenyuAnnotations(new HashMap<>()), "/test", "Exact",
null);
+ ingressReconciler.reconcile(new Request(NAMESPACE, INGRESS_NAME));
+
+ when(ingressIndexer.getByKey(NAMESPACE + "/" +
INGRESS_NAME)).thenReturn(null);
+ Result result = ingressReconciler.reconcile(new Request(NAMESPACE,
INGRESS_NAME));
+
+ assertEquals(new Result(false), result);
+ verify(shenyuCacheRepository, atLeastOnce())
+ .deleteSelectorData(eq(PluginEnum.DIVIDE.getName()),
eq(savedSelectorData.get(0).getId()));
+ assertNull(IngressCache.getInstance().get(NAMESPACE, INGRESS_NAME));
+ assertNull(IngressSelectorCache.getInstance().get(NAMESPACE,
INGRESS_NAME, PluginEnum.DIVIDE.getName()));
+ assertTrue(ServiceIngressCache.getInstance().getIngressName(NAMESPACE,
SERVICE_NAME).isEmpty());
+ }
+
+ private Map<String, String> shenyuAnnotations(final Map<String, String>
annotations) {
+ Map<String, String> allAnnotations = new HashMap<>();
+ allAnnotations.put(IngressConstants.K8S_INGRESS_CLASS_ANNOTATION_KEY,
IngressConstants.SHENYU_INGRESS_CLASS);
+ allAnnotations.putAll(annotations);
+ return allAnnotations;
+ }
+
+ private void mockIngress(final Map<String, String> annotations, final
String path,
+ final String pathType, final String
ingressClassName) {
+ V1IngressRule rule = new
V1IngressRuleBuilder().withNewHttp().withPaths(new V1HTTPIngressPathBuilder()
+ .withPath(path)
+ .withPathType(pathType)
+ .withNewBackend()
+
.withNewService().withName(SERVICE_NAME).withNewPort().withNumber(BACKEND_PORT).endPort().endService()
+ .endBackend()
+ .build())
+ .endHttp().build();
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName(INGRESS_NAME).withNamespace(NAMESPACE)
+ .withAnnotations(annotations).withLabels(new
HashMap<>()).endMetadata()
+
.withNewSpec().withRules(rule).withIngressClassName(ingressClassName).endSpec()
+ .withKind("Ingress")
+ .build();
+ when(ingressIndexer.getByKey(NAMESPACE + "/" +
INGRESS_NAME)).thenReturn(ingress);
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/ContextPathParserTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/ContextPathParserTest.java
new file mode 100644
index 0000000000..778875576d
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/ContextPathParserTest.java
@@ -0,0 +1,202 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.parser;
+
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.informer.cache.Lister;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPath;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1Service;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.dto.convert.rule.impl.ContextMappingRuleHandle;
+import org.apache.shenyu.common.enums.MatchModeEnum;
+import org.apache.shenyu.common.enums.OperatorEnum;
+import org.apache.shenyu.common.enums.ParamTypeEnum;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.k8s.common.IngressConfiguration;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+
+/**
+ * Test for {@link ContextPathParser}.
+ */
+public class ContextPathParserTest {
+
+ private static final String NAMESPACE = "context-path-namespace";
+
+ private Lister<V1Service> serviceLister;
+
+ private Lister<V1Endpoints> endpointsLister;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ serviceLister = new Lister<>(mock(Indexer.class));
+ endpointsLister = new Lister<>(mock(Indexer.class));
+ }
+
+ @Test
+ public void testParseContextPathSelectorAndRule() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_PATH, "/api");
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_ADD_PREFIX,
"/prefix");
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_ADD_PREFIXED,
"true");
+
+ ShenyuMemoryConfig config = parse(createIngress(annotations,
createRule("www.example.com",
+ createPath("/context", "Prefix"))));
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(1, routeConfigList.size());
+
+ SelectorData selectorData = routeConfigList.get(0).getSelectorData();
+ assertEquals(PluginEnum.CONTEXT_PATH.getName(),
selectorData.getPluginName());
+ assertEquals(String.valueOf(PluginEnum.CONTEXT_PATH.getCode()),
selectorData.getPluginId());
+ assertEquals("/context", selectorData.getName());
+ assertEquals(MatchModeEnum.AND.getCode(), selectorData.getMatchMode());
+ assertEquals(2, selectorData.getConditionList().size());
+ assertEquals(ParamTypeEnum.DOMAIN.getName(),
selectorData.getConditionList().get(0).getParamType());
+ assertEquals(OperatorEnum.EQ.getAlias(),
selectorData.getConditionList().get(0).getOperator());
+ assertEquals("www.example.com",
selectorData.getConditionList().get(0).getParamValue());
+ assertEquals(ParamTypeEnum.URI.getName(),
selectorData.getConditionList().get(1).getParamType());
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
selectorData.getConditionList().get(1).getOperator());
+ assertEquals("/context",
selectorData.getConditionList().get(1).getParamValue());
+
+ List<RuleData> ruleDataList = routeConfigList.get(0).getRuleDataList();
+ assertEquals(1, ruleDataList.size());
+ RuleData ruleData = ruleDataList.get(0);
+ assertEquals("/api", ruleData.getName());
+ assertEquals(PluginEnum.CONTEXT_PATH.getName(),
ruleData.getPluginName());
+ assertEquals(MatchModeEnum.AND.getCode(), ruleData.getMatchMode());
+ assertEquals(1, ruleData.getConditionDataList().size());
+ assertEquals(OperatorEnum.PATH_PATTERN.getAlias(),
ruleData.getConditionDataList().get(0).getOperator());
+ assertEquals("/api/**",
ruleData.getConditionDataList().get(0).getParamValue());
+
+ ContextMappingRuleHandle ruleHandle =
GsonUtils.getInstance().fromJson(ruleData.getHandle(),
ContextMappingRuleHandle.class);
+ assertEquals("/api", ruleHandle.getContextPath());
+ assertEquals("/prefix", ruleHandle.getAddPrefix());
+ assertTrue(ruleHandle.getAddPrefixed());
+ }
+
+ @Test
+ public void testParseRuleHandleDefaults() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_PATH, "/api");
+
+ ShenyuMemoryConfig config = parse(createIngress(annotations,
createRule(null, createPath("/context", "Prefix"))));
+
+ List<RuleData> ruleDataList =
config.getRouteConfigList().get(0).getRuleDataList();
+ ContextMappingRuleHandle ruleHandle =
GsonUtils.getInstance().fromJson(ruleDataList.get(0).getHandle(),
ContextMappingRuleHandle.class);
+ assertEquals("/api", ruleHandle.getContextPath());
+ assertNull(ruleHandle.getAddPrefix());
+ assertFalse(ruleHandle.getAddPrefixed());
+ }
+
+ @Test
+ public void testParsePathTypeToOperatorMapping() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_PATH, "/api");
+
+ assertEquals(OperatorEnum.EQ.getAlias(),
parsePathOperator(annotations, "Exact"));
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
parsePathOperator(annotations, "Prefix"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator(annotations, "ImplementationSpecific"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator(annotations, "Unknown"));
+ }
+
+ @Test
+ public void testParseIngressRuleWithNullPath() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_PATH, "/api");
+
+ ShenyuMemoryConfig config = parse(createIngress(annotations,
createRule(null, createPath(null, "Prefix"))));
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(0, routeConfigList.size());
+ }
+
+ @Test
+ public void testParseIngressWithoutRules() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("context-path-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(new HashMap<>()).endMetadata()
+ .withNewSpec().endSpec()
+ .build();
+
+ ShenyuMemoryConfig config = parse(ingress);
+
+ assertNull(config.getRouteConfigList());
+ }
+
+ @Test
+ public void testParseIngressWithoutSpec() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("context-path-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(new HashMap<>()).endMetadata()
+ .build();
+
+ ShenyuMemoryConfig config = parse(ingress);
+
+ assertNull(config.getRouteConfigList());
+ }
+
+ private String parsePathOperator(final Map<String, String> annotations,
final String pathType) {
+ ShenyuMemoryConfig config = parse(createIngress(annotations,
createRule(null, createPath("/context", pathType))));
+ return
config.getRouteConfigList().get(0).getSelectorData().getConditionList().get(0).getOperator();
+ }
+
+ private ShenyuMemoryConfig parse(final V1Ingress ingress) {
+ return new ContextPathParser(serviceLister,
endpointsLister).parse(ingress, null);
+ }
+
+ private V1Ingress createIngress(final Map<String, String> annotations,
final V1IngressRule rule) {
+ return new V1IngressBuilder()
+
.withNewMetadata().withName("context-path-ingress").withNamespace(NAMESPACE).withAnnotations(annotations).endMetadata()
+ .withNewSpec().withRules(rule).endSpec()
+ .withKind("Ingress")
+ .build();
+ }
+
+ private V1IngressRule createRule(final String host, final
V1HTTPIngressPath path) {
+ return new
V1IngressRuleBuilder().withHost(host).withNewHttp().withPaths(path).endHttp().build();
+ }
+
+ private V1HTTPIngressPath createPath(final String path, final String
pathType) {
+ return new
V1HTTPIngressPathBuilder().withPath(path).withPathType(pathType).build();
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/GrpcParserTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/GrpcParserTest.java
new file mode 100644
index 0000000000..0fc98e02ad
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/GrpcParserTest.java
@@ -0,0 +1,275 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.parser;
+
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.informer.cache.Lister;
+import io.kubernetes.client.openapi.models.V1EndpointAddress;
+import io.kubernetes.client.openapi.models.V1EndpointSubsetBuilder;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1EndpointsBuilder;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1Service;
+import io.kubernetes.client.openapi.models.V1ServiceBuilder;
+import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.dto.convert.selector.GrpcUpstream;
+import org.apache.shenyu.common.enums.OperatorEnum;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.k8s.common.IngressConfiguration;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test for {@link GrpcParser}.
+ */
+public class GrpcParserTest {
+
+ private static final String NAMESPACE = "grpc-namespace";
+
+ private static final String SERVICE_NAME = "grpc-service";
+
+ private static final String METADATA_SERVICE_NAME =
"grpc-metadata-service";
+
+ private Indexer<V1Service> serviceIndexer;
+
+ private Indexer<V1Endpoints> endpointsIndexer;
+
+ private Lister<V1Service> serviceLister;
+
+ private Lister<V1Endpoints> endpointsLister;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ serviceIndexer = mock(Indexer.class);
+ endpointsIndexer = mock(Indexer.class);
+ serviceLister = new Lister<>(serviceIndexer);
+ endpointsLister = new Lister<>(endpointsIndexer);
+ }
+
+ @Test
+ public void testParseIngressRuleWithMetadataLabels() {
+ mockEndpoints("10.0.0.1");
+ Map<String, String> metadataAnnotations = new HashMap<>();
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_APP_NAME,
"grpc-app");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_PATH,
"/grpc/hello");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_RPC_TYPE, "grpc");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_SERVICE_NAME,
"hello.HelloService");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_METHOD_NAME,
"hello");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_PARAMS_TYPE,
"hello.HelloRequest");
+ metadataAnnotations.put(IngressConstants.PLUGIN_GRPC_RPC_EXPAND,
"{\"timeout\":5000}");
+ mockService(METADATA_SERVICE_NAME, metadataAnnotations);
+
+ Map<String, String> labels = new HashMap<>();
+ labels.put("shenyu.apache.org/metadata-labels-1",
METADATA_SERVICE_NAME);
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.LOADBALANCER_ANNOTATION_KEY, "hash");
+
+ V1Ingress ingress = createIngress(annotations, labels,
createRule("www.example.com", "/grpc", "Prefix"));
+
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(1, routeConfigList.size());
+
+ SelectorData selectorData = routeConfigList.get(0).getSelectorData();
+ assertEquals(PluginEnum.GRPC.getName(), selectorData.getPluginName());
+ assertEquals(String.valueOf(PluginEnum.GRPC.getCode()),
selectorData.getPluginId());
+ assertEquals("/grpc", selectorData.getName());
+ assertEquals(2, selectorData.getConditionList().size());
+ assertEquals("www.example.com",
selectorData.getConditionList().get(0).getParamValue());
+ assertEquals(OperatorEnum.EQ.getAlias(),
selectorData.getConditionList().get(0).getOperator());
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
selectorData.getConditionList().get(1).getOperator());
+
+ List<GrpcUpstream> upstreams =
GsonUtils.getInstance().fromList(selectorData.getHandle(), GrpcUpstream.class);
+ assertEquals(1, upstreams.size());
+ assertEquals("10.0.0.1:50051", upstreams.get(0).getUpstreamUrl());
+ assertEquals(100, upstreams.get(0).getWeight());
+
+ List<RuleData> ruleDataList = routeConfigList.get(0).getRuleDataList();
+ assertEquals(1, ruleDataList.size());
+ assertEquals("/grpc/hello", ruleDataList.get(0).getName());
+ assertEquals(PluginEnum.GRPC.getName(),
ruleDataList.get(0).getPluginName());
+ assertEquals(1, ruleDataList.get(0).getConditionDataList().size());
+ assertEquals(OperatorEnum.EQ.getAlias(),
ruleDataList.get(0).getConditionDataList().get(0).getOperator());
+ assertEquals("/grpc/hello",
ruleDataList.get(0).getConditionDataList().get(0).getParamValue());
+
+ List<MetaData> metaDataList = routeConfigList.get(0).getMetaDataList();
+ assertEquals(1, metaDataList.size());
+ assertEquals("grpc-app", metaDataList.get(0).getAppName());
+ assertEquals("/grpc/hello", metaDataList.get(0).getPath());
+ assertEquals("grpc", metaDataList.get(0).getRpcType());
+ assertEquals("hello.HelloService",
metaDataList.get(0).getServiceName());
+ assertEquals("hello", metaDataList.get(0).getMethodName());
+ assertEquals("hello.HelloRequest",
metaDataList.get(0).getParameterTypes());
+ assertTrue(metaDataList.get(0).getEnabled());
+ }
+
+ @Test
+ public void testParseGlobalDefaultBackendDefaults() {
+ mockEndpoints("10.0.0.1");
+
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("grpc-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(new HashMap<>()).withLabels(new
HashMap<>()).endMetadata()
+ .withNewSpec()
+
.withNewDefaultBackend().withNewService().withName(SERVICE_NAME).withNewPort().withNumber(50051).endPort().endService().endDefaultBackend()
+ .endSpec()
+ .withKind("Ingress")
+ .build();
+
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNotNull(config.getGlobalDefaultBackend());
+ IngressConfiguration defaultRouteConfig =
config.getGlobalDefaultBackend().getRight();
+ assertEquals("grpc-selector",
defaultRouteConfig.getSelectorData().getName());
+ assertEquals(PluginEnum.GRPC.getName(),
defaultRouteConfig.getSelectorData().getPluginName());
+ assertEquals("grpc-rule",
defaultRouteConfig.getRuleDataList().get(0).getName());
+ assertEquals("/**",
defaultRouteConfig.getRuleDataList().get(0).getConditionDataList().get(0).getParamValue());
+ assertEquals(OperatorEnum.PATH_PATTERN.getAlias(),
defaultRouteConfig.getRuleDataList().get(0).getConditionDataList().get(0).getOperator());
+
+ List<GrpcUpstream> upstreams = GsonUtils.getInstance()
+ .fromList(defaultRouteConfig.getSelectorData().getHandle(),
GrpcUpstream.class);
+ assertEquals(1, upstreams.size());
+ assertEquals("10.0.0.1:50051", upstreams.get(0).getUpstreamUrl());
+ assertEquals(50, upstreams.get(0).getWeight());
+
+ MetaData metaData = defaultRouteConfig.getMetaDataList().get(0);
+ assertEquals("grpc", metaData.getAppName());
+ assertEquals("/grpc/helloService/hello", metaData.getPath());
+ assertEquals("grpc", metaData.getRpcType());
+ assertEquals("hello.HelloService", metaData.getServiceName());
+ assertEquals("/grpc", metaData.getContextPath());
+ assertTrue(metaData.getEnabled());
+ }
+
+ @Test
+ public void testParseIngressRuleWithoutMetadataLabels() {
+ mockEndpoints("10.0.0.1");
+
+ V1Ingress ingress = createIngress(new HashMap<>(), new HashMap<>(),
createRule(null, "/grpc", "Prefix"));
+
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ IngressConfiguration routeConfig = config.getRouteConfigList().get(0);
+ assertEquals(0, routeConfig.getRuleDataList().size());
+ assertEquals(0, routeConfig.getMetaDataList().size());
+ List<GrpcUpstream> upstreams =
GsonUtils.getInstance().fromList(routeConfig.getSelectorData().getHandle(),
GrpcUpstream.class);
+ assertEquals("10.0.0.1:50051", upstreams.get(0).getUpstreamUrl());
+ }
+
+ @Test
+ public void testParseIngressRuleWithoutEndpointsSubsets() {
+ when(endpointsIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(new V1EndpointsBuilder()
+
.withNewMetadata().withNamespace(NAMESPACE).withName(SERVICE_NAME).endMetadata()
+ .build());
+
+ V1Ingress ingress = createIngress(new HashMap<>(), new HashMap<>(),
createRule(null, "/grpc", "Prefix"));
+
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertEquals("[]",
config.getRouteConfigList().get(0).getSelectorData().getHandle());
+ }
+
+ @Test
+ public void testParsePathTypeToOperatorMapping() {
+ mockEndpoints("10.0.0.1");
+
+ assertEquals(OperatorEnum.EQ.getAlias(), parsePathOperator("Exact"));
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
parsePathOperator("Prefix"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator("ImplementationSpecific"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator("Unknown"));
+ }
+
+ @Test
+ public void testParseIngressWithoutSpec() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("grpc-ingress").withNamespace(NAMESPACE).endMetadata()
+ .build();
+
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNull(config.getRouteConfigList());
+ assertNull(config.getGlobalDefaultBackend());
+ }
+
+ private String parsePathOperator(final String pathType) {
+ V1Ingress ingress = createIngress(new HashMap<>(), new HashMap<>(),
createRule(null, "/grpc", pathType));
+ ShenyuMemoryConfig config = new GrpcParser(serviceLister,
endpointsLister).parse(ingress, null);
+ return
config.getRouteConfigList().get(0).getSelectorData().getConditionList().get(0).getOperator();
+ }
+
+ private void mockEndpoints(final String... ips) {
+ V1EndpointAddress[] addresses = new V1EndpointAddress[ips.length];
+ for (int i = 0; i < ips.length; i++) {
+ addresses[i] = new V1EndpointAddress().ip(ips[i]);
+ }
+ V1Endpoints endpoints = new V1EndpointsBuilder()
+
.withNewMetadata().withNamespace(NAMESPACE).withName(SERVICE_NAME).endMetadata()
+ .withSubsets(new
V1EndpointSubsetBuilder().withAddresses(addresses).build())
+ .build();
+ when(endpointsIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(endpoints);
+ }
+
+ private void mockService(final String serviceName, final Map<String,
String> annotations) {
+ when(serviceIndexer.getByKey(NAMESPACE + "/" +
serviceName)).thenReturn(new V1ServiceBuilder()
+
.withNewMetadata().withName(serviceName).withNamespace(NAMESPACE).withAnnotations(annotations).endMetadata()
+ .withKind("Service").build());
+ }
+
+ private V1Ingress createIngress(final Map<String, String> annotations,
final Map<String, String> labels, final V1IngressRule rule) {
+ return new V1IngressBuilder()
+
.withNewMetadata().withName("grpc-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(annotations).withLabels(labels).endMetadata()
+ .withNewSpec().withRules(rule).endSpec()
+ .withKind("Ingress")
+ .build();
+ }
+
+ private V1IngressRule createRule(final String host, final String path,
final String pathType) {
+ return new
V1IngressRuleBuilder().withHost(host).withNewHttp().withPaths(new
V1HTTPIngressPathBuilder()
+ .withPath(path)
+ .withPathType(pathType)
+ .withNewBackend()
+
.withNewService().withName(SERVICE_NAME).withNewPort().withNumber(50051).endPort().endService()
+ .endBackend()
+ .build())
+ .endHttp().build();
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/IngressParserTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/IngressParserTest.java
new file mode 100644
index 0000000000..707a087359
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/IngressParserTest.java
@@ -0,0 +1,195 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.parser;
+
+import io.kubernetes.client.informer.SharedIndexInformer;
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.openapi.apis.CoreV1Api;
+import io.kubernetes.client.openapi.models.V1EndpointAddress;
+import io.kubernetes.client.openapi.models.V1EndpointSubsetBuilder;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1EndpointsBuilder;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1Service;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.k8s.common.IngressConfiguration;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test for {@link IngressParser}.
+ */
+public class IngressParserTest {
+
+ private static final String NAMESPACE = "dispatch-namespace";
+
+ private static final String SERVICE_NAME = "dispatch-service";
+
+ private IngressParser ingressParser;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ SharedIndexInformer<V1Service> serviceInformer =
mock(SharedIndexInformer.class);
+ when(serviceInformer.getIndexer()).thenReturn(mock(Indexer.class));
+
+ Indexer<V1Endpoints> endpointsIndexer = mock(Indexer.class);
+ SharedIndexInformer<V1Endpoints> endpointsInformer =
mock(SharedIndexInformer.class);
+ when(endpointsInformer.getIndexer()).thenReturn(endpointsIndexer);
+
+ V1Endpoints endpoints = new V1EndpointsBuilder()
+
.withNewMetadata().withNamespace(NAMESPACE).withName(SERVICE_NAME).endMetadata()
+ .withSubsets(new V1EndpointSubsetBuilder().withAddresses(new
V1EndpointAddress().ip("10.0.0.1")).build())
+ .build();
+ when(endpointsIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(endpoints);
+
+ ingressParser = new IngressParser(serviceInformer, endpointsInformer);
+ }
+
+ @Test
+ public void testDispatchToDivideWhenNoPluginAnnotationEnabled() {
+ List<ShenyuMemoryConfig> configs = parse(new HashMap<>(), "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.DIVIDE.getName()),
pluginNames(configs));
+ }
+
+ @Test
+ public void testDispatchAddsContextPathConfigForPrefixPath() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_CONTEXT_PATH_PATH, "/test");
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Prefix");
+
+ assertEquals(2, configs.size());
+ List<String> pluginNames = pluginNames(configs);
+ assertTrue(pluginNames.contains(PluginEnum.CONTEXT_PATH.getName()));
+ assertTrue(pluginNames.contains(PluginEnum.DIVIDE.getName()));
+ }
+
+ @Test
+ public void testDispatchToDubboWhenEnabled() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_DUBBO_ENABLED, "true");
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.DUBBO.getName()),
pluginNames(configs));
+ }
+
+ @Test
+ public void testDispatchToWebSocketWhenEnabled() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_WEB_SOCKET_ENABLED, "true");
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.WEB_SOCKET.getName()),
pluginNames(configs));
+ }
+
+ @Test
+ public void testDispatchToGrpcWhenEnabled() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_GRPC_ENABLED, "true");
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.GRPC.getName()),
pluginNames(configs));
+ }
+
+ @Test
+ public void testDispatchToSofaWhenEnabled() {
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.PLUGIN_SOFA_ENABLED, "true");
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.SOFA.getName()),
pluginNames(configs));
+ }
+
+ @Test
+ public void testDispatchToDivideWhenAllPluginAnnotationsDisabled() {
+ Map<String, String> annotations = new HashMap<>();
+ for (String key : Arrays.asList(IngressConstants.PLUGIN_DUBBO_ENABLED,
IngressConstants.PLUGIN_WEB_SOCKET_ENABLED,
+ IngressConstants.PLUGIN_BRPC_ENABLED,
IngressConstants.PLUGIN_GRPC_ENABLED, IngressConstants.PLUGIN_SOFA_ENABLED)) {
+ annotations.put(key, "false");
+ }
+
+ List<ShenyuMemoryConfig> configs = parse(annotations, "Exact");
+
+ assertEquals(1, configs.size());
+ assertEquals(Arrays.asList(PluginEnum.DIVIDE.getName()),
pluginNames(configs));
+ }
+
+ private List<ShenyuMemoryConfig> parse(final Map<String, String>
annotations, final String pathType) {
+ Map<String, String> allAnnotations = new HashMap<>();
+ allAnnotations.put(IngressConstants.K8S_INGRESS_CLASS_ANNOTATION_KEY,
IngressConstants.SHENYU_INGRESS_CLASS);
+ allAnnotations.putAll(annotations);
+
+ V1IngressRule rule = new
V1IngressRuleBuilder().withNewHttp().withPaths(new V1HTTPIngressPathBuilder()
+ .withPath("/test")
+ .withPathType(pathType)
+ .withNewBackend()
+
.withNewService().withName(SERVICE_NAME).withNewPort().withNumber(9090).endPort().endService()
+ .endBackend()
+ .build())
+ .endHttp().build();
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("dispatch-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(allAnnotations).withLabels(new
HashMap<>()).endMetadata()
+ .withNewSpec().withRules(rule).endSpec()
+ .withKind("Ingress")
+ .build();
+
+ return ingressParser.parse(ingress, mock(CoreV1Api.class));
+ }
+
+ private List<String> pluginNames(final List<ShenyuMemoryConfig> configs) {
+ return configs.stream()
+ .map(ShenyuMemoryConfig::getRouteConfigList)
+ .filter(Objects::nonNull)
+ .flatMap(List::stream)
+ .map(IngressConfiguration::getSelectorData)
+ .filter(Objects::nonNull)
+ .map(SelectorData::getPluginName)
+ .collect(Collectors.toList());
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/SofaParserTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/SofaParserTest.java
new file mode 100644
index 0000000000..f4170bd763
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/SofaParserTest.java
@@ -0,0 +1,233 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.parser;
+
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.informer.cache.Lister;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1Service;
+import io.kubernetes.client.openapi.models.V1ServiceBuilder;
+import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.OperatorEnum;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.k8s.common.IngressConfiguration;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test for {@link SofaParser}.
+ */
+public class SofaParserTest {
+
+ private static final String NAMESPACE = "sofa-namespace";
+
+ private static final String SERVICE_NAME = "sofa-service";
+
+ private static final String METADATA_SERVICE_NAME =
"sofa-metadata-service";
+
+ private Indexer<V1Service> serviceIndexer;
+
+ private Indexer<V1Endpoints> endpointsIndexer;
+
+ private Lister<V1Service> serviceLister;
+
+ private Lister<V1Endpoints> endpointsLister;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ serviceIndexer = mock(Indexer.class);
+ endpointsIndexer = mock(Indexer.class);
+ serviceLister = new Lister<>(serviceIndexer);
+ endpointsLister = new Lister<>(endpointsIndexer);
+ }
+
+ @Test
+ public void testParseIngressRuleWithMetadataLabels() {
+ Map<String, String> metadataAnnotations = new HashMap<>();
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_APP_NAME,
"sofa-app");
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_PATH,
"/sofa/findById");
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_RPC_TYPE, "sofa");
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_SERVICE_NAME,
"org.apache.shenyu.examples.sofa.api.SofaTestService");
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_METHOD_NAME,
"findById");
+ metadataAnnotations.put(IngressConstants.PLUGIN_SOFA_PARAMS_TYPE,
"java.lang.String");
+ mockService(METADATA_SERVICE_NAME, metadataAnnotations);
+
+ Map<String, String> labels = new HashMap<>();
+ labels.put("shenyu.apache.org/metadata-labels-1",
METADATA_SERVICE_NAME);
+ Map<String, String> annotations = new HashMap<>();
+ annotations.put(IngressConstants.LOADBALANCER_ANNOTATION_KEY, "hash");
+ annotations.put(IngressConstants.RETRY_ANNOTATION_KEY, "2");
+
+ V1Ingress ingress = createIngress(annotations, labels,
createRule("www.example.com", "/sofa", "Prefix"));
+
+ ShenyuMemoryConfig config = new SofaParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(1, routeConfigList.size());
+
+ SelectorData selectorData = routeConfigList.get(0).getSelectorData();
+ assertEquals(PluginEnum.SOFA.getName(), selectorData.getPluginName());
+ assertEquals(String.valueOf(PluginEnum.SOFA.getCode()),
selectorData.getPluginId());
+ assertEquals("/sofa", selectorData.getName());
+ assertEquals(2, selectorData.getConditionList().size());
+ assertEquals("www.example.com",
selectorData.getConditionList().get(0).getParamValue());
+ assertEquals(OperatorEnum.EQ.getAlias(),
selectorData.getConditionList().get(0).getOperator());
+ assertEquals("/sofa",
selectorData.getConditionList().get(1).getParamValue());
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
selectorData.getConditionList().get(1).getOperator());
+
+ List<RuleData> ruleDataList = routeConfigList.get(0).getRuleDataList();
+ assertEquals(1, ruleDataList.size());
+ assertEquals("/sofa/findById", ruleDataList.get(0).getName());
+ assertEquals(PluginEnum.SOFA.getName(),
ruleDataList.get(0).getPluginName());
+ assertEquals(1, ruleDataList.get(0).getConditionDataList().size());
+ assertEquals(OperatorEnum.EQ.getAlias(),
ruleDataList.get(0).getConditionDataList().get(0).getOperator());
+ assertEquals("/sofa/findById",
ruleDataList.get(0).getConditionDataList().get(0).getParamValue());
+
+ List<MetaData> metaDataList = routeConfigList.get(0).getMetaDataList();
+ assertEquals(1, metaDataList.size());
+ assertEquals("sofa-app", metaDataList.get(0).getAppName());
+ assertEquals("/sofa/findById", metaDataList.get(0).getPath());
+ assertEquals("sofa", metaDataList.get(0).getRpcType());
+ assertEquals("org.apache.shenyu.examples.sofa.api.SofaTestService",
metaDataList.get(0).getServiceName());
+ assertEquals("findById", metaDataList.get(0).getMethodName());
+ assertEquals("java.lang.String",
metaDataList.get(0).getParameterTypes());
+ assertTrue(metaDataList.get(0).getEnabled());
+ }
+
+ @Test
+ public void testParseGlobalDefaultBackendDefaults() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("sofa-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(new HashMap<>()).withLabels(new
HashMap<>()).endMetadata()
+ .withNewSpec()
+
.withNewDefaultBackend().withNewService().withName(SERVICE_NAME).withNewPort().withNumber(12200).endPort().endService().endDefaultBackend()
+ .endSpec()
+ .withKind("Ingress")
+ .build();
+
+ ShenyuMemoryConfig config = new SofaParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNotNull(config.getGlobalDefaultBackend());
+ IngressConfiguration defaultRouteConfig =
config.getGlobalDefaultBackend().getRight();
+ assertEquals("sofa-selector",
defaultRouteConfig.getSelectorData().getName());
+ assertEquals(PluginEnum.SOFA.getName(),
defaultRouteConfig.getSelectorData().getPluginName());
+ assertEquals(IngressConstants.ID,
defaultRouteConfig.getSelectorData().getId());
+ assertEquals("sofa-rule",
defaultRouteConfig.getRuleDataList().get(0).getName());
+ assertEquals(IngressConstants.ID,
defaultRouteConfig.getRuleDataList().get(0).getSelectorId());
+ assertEquals("/**",
defaultRouteConfig.getRuleDataList().get(0).getConditionDataList().get(0).getParamValue());
+ assertEquals(OperatorEnum.PATH_PATTERN.getAlias(),
defaultRouteConfig.getRuleDataList().get(0).getConditionDataList().get(0).getOperator());
+
+ MetaData metaData = defaultRouteConfig.getMetaDataList().get(0);
+ assertEquals("sofa", metaData.getAppName());
+ assertEquals("/sofa/findAll", metaData.getPath());
+ assertEquals("sofa", metaData.getRpcType());
+ assertEquals("findAll", metaData.getServiceName());
+ assertEquals("methodName", metaData.getMethodName());
+ assertEquals("/sofa", metaData.getContextPath());
+ assertEquals("", metaData.getParameterTypes());
+ assertTrue(metaData.getEnabled());
+ }
+
+ @Test
+ public void testParseIngressRuleWithoutMetadataLabels() {
+ V1Ingress ingress = createIngress(new HashMap<>(), new HashMap<>(),
createRule(null, "/sofa", "Prefix"));
+
+ ShenyuMemoryConfig config = new SofaParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ IngressConfiguration routeConfig = config.getRouteConfigList().get(0);
+ assertEquals(0, routeConfig.getRuleDataList().size());
+ assertEquals(0, routeConfig.getMetaDataList().size());
+ }
+
+ @Test
+ public void testParsePathTypeToOperatorMapping() {
+ assertEquals(OperatorEnum.EQ.getAlias(), parsePathOperator("Exact"));
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
parsePathOperator("Prefix"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator("ImplementationSpecific"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parsePathOperator("Unknown"));
+ }
+
+ @Test
+ public void testParseIngressWithoutRulesAndDefaultBackend() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("sofa-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(new HashMap<>()).withLabels(new
HashMap<>()).endMetadata()
+ .withNewSpec().endSpec()
+ .withKind("Ingress")
+ .build();
+
+ ShenyuMemoryConfig config = new SofaParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNull(config.getRouteConfigList());
+ assertNull(config.getGlobalDefaultBackend());
+ }
+
+ private String parsePathOperator(final String pathType) {
+ V1Ingress ingress = createIngress(new HashMap<>(), new HashMap<>(),
createRule(null, "/sofa", pathType));
+ ShenyuMemoryConfig config = new SofaParser(serviceLister,
endpointsLister).parse(ingress, null);
+ return
config.getRouteConfigList().get(0).getSelectorData().getConditionList().get(0).getOperator();
+ }
+
+ private void mockService(final String serviceName, final Map<String,
String> annotations) {
+ when(serviceIndexer.getByKey(NAMESPACE + "/" +
serviceName)).thenReturn(new V1ServiceBuilder()
+
.withNewMetadata().withName(serviceName).withNamespace(NAMESPACE).withAnnotations(annotations).endMetadata()
+ .withKind("Service").build());
+ }
+
+ private V1Ingress createIngress(final Map<String, String> annotations,
final Map<String, String> labels, final V1IngressRule rule) {
+ return new V1IngressBuilder()
+
.withNewMetadata().withName("sofa-ingress").withNamespace(NAMESPACE)
+ .withAnnotations(annotations).withLabels(labels).endMetadata()
+ .withNewSpec().withRules(rule).endSpec()
+ .withKind("Ingress")
+ .build();
+ }
+
+ private V1IngressRule createRule(final String host, final String path,
final String pathType) {
+ return new
V1IngressRuleBuilder().withHost(host).withNewHttp().withPaths(new
V1HTTPIngressPathBuilder()
+ .withPath(path)
+ .withPathType(pathType)
+ .withNewBackend()
+
.withNewService().withName(SERVICE_NAME).withNewPort().withNumber(12200).endPort().endService()
+ .endBackend()
+ .build())
+ .endHttp().build();
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/WebSocketParserTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/WebSocketParserTest.java
new file mode 100644
index 0000000000..28a3775a55
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/parser/WebSocketParserTest.java
@@ -0,0 +1,310 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.parser;
+
+import io.kubernetes.client.informer.cache.Indexer;
+import io.kubernetes.client.informer.cache.Lister;
+import io.kubernetes.client.openapi.apis.CoreV1Api;
+import io.kubernetes.client.openapi.models.V1EndpointAddress;
+import io.kubernetes.client.openapi.models.V1EndpointSubsetBuilder;
+import io.kubernetes.client.openapi.models.V1Endpoints;
+import io.kubernetes.client.openapi.models.V1EndpointsBuilder;
+import io.kubernetes.client.openapi.models.V1HTTPIngressPathBuilder;
+import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1IngressBuilder;
+import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressRuleBuilder;
+import io.kubernetes.client.openapi.models.V1IngressTLS;
+import io.kubernetes.client.openapi.models.V1IngressTLSBuilder;
+import io.kubernetes.client.openapi.models.V1Secret;
+import io.kubernetes.client.openapi.models.V1SecretBuilder;
+import io.kubernetes.client.openapi.models.V1Service;
+import io.kubernetes.client.openapi.models.V1ServiceBuilder;
+import org.apache.shenyu.common.dto.ConditionData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.dto.convert.selector.WebSocketUpstream;
+import org.apache.shenyu.common.enums.OperatorEnum;
+import org.apache.shenyu.common.enums.ParamTypeEnum;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.k8s.common.IngressConfiguration;
+import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test for {@link WebSocketParser}.
+ */
+public class WebSocketParserTest {
+
+ private static final String NAMESPACE = "ws-namespace";
+
+ private static final String SERVICE_NAME = "ws-service";
+
+ private Indexer<V1Service> serviceIndexer;
+
+ private Indexer<V1Endpoints> endpointsIndexer;
+
+ private Lister<V1Service> serviceLister;
+
+ private Lister<V1Endpoints> endpointsLister;
+
+ @BeforeEach
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ serviceIndexer = mock(Indexer.class);
+ endpointsIndexer = mock(Indexer.class);
+ serviceLister = new Lister<>(serviceIndexer);
+ endpointsLister = new Lister<>(endpointsIndexer);
+ }
+
+ @Test
+ public void testParseIngressRuleToWebSocketSelectorAndRule() {
+ mockEndpoints("10.0.0.1");
+
+ V1Ingress ingress = createIngress(createRule("www.example.com", "/ws",
"Prefix", SERVICE_NAME, 8001),
+ Collections.emptyMap(), null);
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(1, routeConfigList.size());
+
+ SelectorData selectorData = routeConfigList.get(0).getSelectorData();
+ assertNotNull(selectorData);
+ assertEquals(PluginEnum.WEB_SOCKET.getName(),
selectorData.getPluginName());
+ assertEquals(String.valueOf(PluginEnum.WEB_SOCKET.getCode()),
selectorData.getPluginId());
+ assertEquals("/ws", selectorData.getName());
+ assertEquals(2, selectorData.getConditionList().size());
+ assertEquals("www.example.com",
selectorData.getConditionList().get(0).getParamValue());
+ assertEquals(ParamTypeEnum.DOMAIN.getName(),
selectorData.getConditionList().get(0).getParamType());
+ assertEquals(OperatorEnum.EQ.getAlias(),
selectorData.getConditionList().get(0).getOperator());
+ assertEquals("/ws",
selectorData.getConditionList().get(1).getParamValue());
+ assertEquals(ParamTypeEnum.URI.getName(),
selectorData.getConditionList().get(1).getParamType());
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
selectorData.getConditionList().get(1).getOperator());
+
+ List<WebSocketUpstream> upstreams =
GsonUtils.getInstance().fromList(selectorData.getHandle(),
WebSocketUpstream.class);
+ assertEquals(1, upstreams.size());
+ assertEquals("10.0.0.1:8001", upstreams.get(0).getUpstreamUrl());
+ assertEquals("ws://", upstreams.get(0).getProtocol());
+ assertEquals(100, upstreams.get(0).getWeight());
+
+ List<RuleData> ruleDataList = routeConfigList.get(0).getRuleDataList();
+ assertEquals(1, ruleDataList.size());
+ assertEquals("/ws", ruleDataList.get(0).getName());
+ assertEquals(PluginEnum.WEB_SOCKET.getName(),
ruleDataList.get(0).getPluginName());
+ assertThat(ruleDataList.get(0).getHandle(),
containsString("\"timeout\":3000"));
+ }
+
+ @Test
+ public void testPathTypeToOperatorMapping() {
+ mockEndpoints("10.0.0.1");
+
+ assertEquals(OperatorEnum.EQ.getAlias(),
parseFirstPathConditionOperator("Exact"));
+ assertEquals(OperatorEnum.STARTS_WITH.getAlias(),
parseFirstPathConditionOperator("Prefix"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parseFirstPathConditionOperator("ImplementationSpecific"));
+ assertEquals(OperatorEnum.MATCH.getAlias(),
parseFirstPathConditionOperator("Unknown"));
+ }
+
+ @Test
+ public void testParseIngressRuleWithoutHost() {
+ mockEndpoints("10.0.0.1");
+
+ V1Ingress ingress = createIngress(createRule(null, "/ws", "Prefix",
SERVICE_NAME, 8001),
+ Collections.emptyMap(), null);
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<ConditionData> conditionList =
config.getRouteConfigList().get(0).getSelectorData().getConditionList();
+ assertEquals(1, conditionList.size());
+ assertEquals(ParamTypeEnum.URI.getName(),
conditionList.get(0).getParamType());
+ }
+
+ @Test
+ public void testParseIngressRuleSkipsPathWithoutPath() {
+ V1IngressRule rule = new V1IngressRuleBuilder()
+ .withNewHttp().withPaths(new
V1HTTPIngressPathBuilder().withNewBackend().withNewService()
+
.withName(SERVICE_NAME).withNewPort().withNumber(8001).endPort().endService().endBackend().build())
+ .endHttp().build();
+ V1Ingress ingress = createIngress(rule, Collections.emptyMap(), null);
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<IngressConfiguration> routeConfigList =
config.getRouteConfigList();
+ assertNotNull(routeConfigList);
+ assertEquals(0, routeConfigList.size());
+ }
+
+ @Test
+ public void testParseGlobalDefaultBackendWithServiceProtocolAnnotation() {
+ mockEndpoints("10.0.0.1", "10.0.0.2");
+ Map<String, String> serviceAnnotations = new HashMap<>();
+
serviceAnnotations.put(IngressConstants.UPSTREAMS_PROTOCOL_ANNOTATION_KEY,
"ws://,wss://");
+ mockService(serviceAnnotations);
+
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("ws-ingress").withNamespace(NAMESPACE).withAnnotations(Collections.emptyMap()).endMetadata()
+ .withNewSpec()
+
.withNewDefaultBackend().withNewService().withName(SERVICE_NAME).withNewPort().withNumber(8001).endPort().endService().endDefaultBackend()
+ .endSpec()
+ .build();
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNotNull(config.getGlobalDefaultBackend());
+ SelectorData selectorData =
config.getGlobalDefaultBackend().getRight().getSelectorData();
+ assertEquals(PluginEnum.WEB_SOCKET.getName(),
selectorData.getPluginName());
+ assertEquals("default-selector", selectorData.getName());
+
+ List<WebSocketUpstream> upstreams =
GsonUtils.getInstance().fromList(selectorData.getHandle(),
WebSocketUpstream.class);
+ assertEquals(2, upstreams.size());
+ assertEquals("10.0.0.1:8001", upstreams.get(0).getUpstreamUrl());
+ assertEquals("ws://", upstreams.get(0).getProtocol());
+ assertEquals("10.0.0.2:8001", upstreams.get(1).getUpstreamUrl());
+ assertEquals("wss://", upstreams.get(1).getProtocol());
+ }
+
+ @Test
+ public void
testParseIngressRuleFallsBackToEmptyUpstreamWhenEndpointsMissing() {
+ V1Ingress ingress = createIngress(createRule(null, "/ws", "Prefix",
SERVICE_NAME, 8001),
+ Collections.emptyMap(), null);
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertEquals("[]",
config.getRouteConfigList().get(0).getSelectorData().getHandle());
+ }
+
+ @Test
+ public void testParseIngressRuleUsesDefaultUpstreamFromDefaultBackend() {
+ mockEndpoints("10.0.0.1");
+ mockService(Collections.emptyMap());
+
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("ws-ingress").withNamespace(NAMESPACE).withAnnotations(Collections.emptyMap()).endMetadata()
+ .withNewSpec()
+
.withNewDefaultBackend().withNewService().withName(SERVICE_NAME).withNewPort().withNumber(8001).endPort().endService().endDefaultBackend()
+ .withRules(createRule(null, "/ws", "Prefix",
"unknown-service", 8002))
+ .endSpec()
+ .build();
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ List<WebSocketUpstream> upstreams = GsonUtils.getInstance()
+
.fromList(config.getRouteConfigList().get(0).getSelectorData().getHandle(),
WebSocketUpstream.class);
+ assertEquals(1, upstreams.size());
+ assertEquals("10.0.0.1:8001", upstreams.get(0).getUpstreamUrl());
+ }
+
+ @Test
+ public void testParseTlsConfigurations() throws Exception {
+ mockEndpoints("10.0.0.1");
+ CoreV1Api coreV1Api = mock(CoreV1Api.class);
+ Map<String, byte[]> secretData = new HashMap<>();
+ secretData.put("tls.crt", "crt".getBytes());
+ secretData.put("tls.key", "key".getBytes());
+ V1Secret secret = new V1SecretBuilder()
+
.withNewMetadata().withName("ws-tls-secret").withNamespace(NAMESPACE).endMetadata()
+ .withData(secretData)
+ .build();
+ when(coreV1Api.readNamespacedSecret(eq("ws-tls-secret"),
eq(NAMESPACE), anyString())).thenReturn(secret);
+
+ V1Ingress ingress = createIngress(createRule(null, "/ws", "Prefix",
SERVICE_NAME, 8001),
+ Collections.emptyMap(),
+ new
V1IngressTLSBuilder().withHosts(Arrays.asList("www.example.com")).withSecretName("ws-tls-secret").build());
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, coreV1Api);
+
+ assertNotNull(config.getTlsConfigList());
+ assertEquals(1, config.getTlsConfigList().size());
+ assertEquals("www.example.com",
config.getTlsConfigList().get(0).getDomain());
+ }
+
+ @Test
+ public void testParseWithoutSpecReturnsEmptyConfig() {
+ V1Ingress ingress = new V1IngressBuilder()
+
.withNewMetadata().withName("ws-ingress").withNamespace(NAMESPACE).endMetadata()
+ .build();
+
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+
+ assertNull(config.getRouteConfigList());
+ assertNull(config.getGlobalDefaultBackend());
+ }
+
+ private String parseFirstPathConditionOperator(final String pathType) {
+ V1Ingress ingress = createIngress(createRule(null, "/ws", pathType,
SERVICE_NAME, 8001),
+ Collections.emptyMap(), null);
+ ShenyuMemoryConfig config = new WebSocketParser(serviceLister,
endpointsLister).parse(ingress, null);
+ return
config.getRouteConfigList().get(0).getSelectorData().getConditionList().get(0).getOperator();
+ }
+
+ private void mockEndpoints(final String... ips) {
+ V1EndpointAddress[] addresses = Arrays.stream(ips).map(ip -> new
V1EndpointAddress().ip(ip)).toArray(V1EndpointAddress[]::new);
+ V1Endpoints endpoints = new V1EndpointsBuilder()
+
.withNewMetadata().withNamespace(NAMESPACE).withName(SERVICE_NAME).endMetadata()
+ .withSubsets(new
V1EndpointSubsetBuilder().withAddresses(addresses).build())
+ .build();
+ when(endpointsIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(endpoints);
+ }
+
+ private void mockService(final Map<String, String> annotations) {
+ when(serviceIndexer.getByKey(NAMESPACE + "/" +
SERVICE_NAME)).thenReturn(new V1ServiceBuilder()
+
.withNewMetadata().withName(SERVICE_NAME).withNamespace(NAMESPACE).withAnnotations(annotations).endMetadata()
+ .withKind("Service").build());
+ }
+
+ private V1IngressRule createRule(final String host, final String path,
final String pathType,
+ final String serviceName, final int port)
{
+ V1HTTPIngressPathBuilder ingressPathBuilder = new
V1HTTPIngressPathBuilder()
+ .withPath(path)
+ .withPathType(pathType)
+ .withNewBackend()
+
.withNewService().withName(serviceName).withNewPort().withNumber(port).endPort().endService()
+ .endBackend();
+ return new
V1IngressRuleBuilder().withHost(host).withNewHttp().withPaths(ingressPathBuilder.build()).endHttp().build();
+ }
+
+ private V1Ingress createIngress(final V1IngressRule rule, final
Map<String, String> annotations, final V1IngressTLS tls) {
+ List<V1IngressTLS> tlsList = Objects.isNull(tls) ?
Collections.emptyList() : Collections.singletonList(tls);
+ return new V1IngressBuilder()
+
.withNewMetadata().withName("ws-ingress").withNamespace(NAMESPACE).withAnnotations(annotations).endMetadata()
+ .withNewSpec().withRules(rule).withTls(tlsList).endSpec()
+ .withKind("Ingress")
+ .build();
+ }
+}