This is an automated email from the ASF dual-hosted git repository.
liubao pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/servicecomb-java-chassis.git
The following commit(s) were added to refs/heads/master by this push:
new 23e5e4152 [SCB-2886]support web-socket program model (#4415)
23e5e4152 is described below
commit 23e5e4152ed64d19a8b3080e2d25194424ba1474
Author: liubao68 <[email protected]>
AuthorDate: Wed Jul 17 17:21:35 2024 +0800
[SCB-2886]support web-socket program model (#4415)
---
.../common/rest/CommonRestConfiguration.java | 6 +
.../rest/EdgeServerWebSocketInvocationCreator.java | 35 ++---
.../ProviderServerWebSocketInvocationCreator.java | 23 +++-
.../rest/ServerWebSocketInvocationCreator.java | 136 +++++++++++++++++++
.../common/rest/WebSocketTransportContext.java | 22 +--
.../filter/inner/WebSocketServerCodecFilter.java | 124 +++++++++++++++++
.../rest/route}/URLMappedConfigurationItem.java | 2 +-
.../rest/route}/URLMappedConfigurationLoader.java | 2 +-
.../servicecomb/common/rest/route}/Utils.java | 2 +-
.../org/apache/servicecomb/core/CoreConst.java | 2 +
.../org/apache/servicecomb/core/Transport.java | 1 -
.../core/invocation/InvocationCreator.java | 3 -
.../samples/ClientWebsocketController.java | 54 ++++++++
.../consumer/src/main/resources/application.yml | 4 +-
.../gateway/src/main/resources/application.yml | 11 +-
.../servicecomb/samples/WebsocketController.java | 57 ++++++++
.../provider/src/main/resources/application.yml | 4 +-
.../servicecomb/samples/ThirdSvcConfiguration.java | 20 ++-
.../apache/servicecomb/samples/WebsocketIT.java | 57 ++++++++
.../edge/core/CommonHttpEdgeDispatcher.java | 3 +
.../edge/core/DefaultEdgeDispatcher.java | 1 +
.../servicecomb/edge/core/EdgeBootListener.java | 13 --
.../edge/core/EdgeInvocationCreator.java | 1 +
.../edge/core/URLMappedEdgeDispatcher.java | 9 +-
.../edge/core/TestURLMappedEdgeDispatcher.java | 1 +
.../apache/servicecomb/edge/core/TestUtils.java | 1 +
.../foundation/common/net/URIEndpointObject.java | 16 +++
.../discovery/AbstractEndpointDiscoveryFilter.java | 22 ++-
.../foundation/vertx/VertxTLSBuilder.java | 20 +++
.../vertx/client/http/HttpClientOptionsSPI.java | 38 ++++--
.../foundation/vertx/client/http/HttpClients.java | 6 +
.../servicecomb/loadbalance/LoadBalanceFilter.java | 5 +-
.../ProducerServerWebSocketArgMapperFactory.java | 27 ++--
.../rest/common/ProducerServerWebSocketMapper.java | 39 ++----
...s.producer.ProducerContextArgumentMapperFactory | 3 +-
.../swagger/generator/OperationGenerator.java | 6 +
.../swagger/generator/ParameterGenerator.java | 2 +-
.../swagger/generator/SwaggerConst.java | 4 +
.../generator/core/AbstractOperationGenerator.java | 7 +
.../generator/core/ParameterGeneratorContext.java | 22 ++-
.../generator/core/SwaggerGeneratorContext.java | 5 +-
.../parameter/ServerWebSocketContextRegister.java | 19 +--
...cecomb.swagger.generator.SwaggerContextRegister | 3 +-
.../client/TransportRestClientConfiguration.java | 5 +
.../rest/client/WebSocketClientCodecFilter.java | 148 +++++++++++++++++++++
.../client/WebSocketClientTransportContext.java | 23 ++--
.../transport/rest/vertx/RestServerVerticle.java | 15 ++-
.../transport/rest/vertx/WebSocketDispatcher.java | 122 +++++++++++++++++
.../vertx/WebSocketProducerInvocationFlow.java | 51 +++++++
.../transport/rest/vertx/WebSocketTransport.java | 26 ++--
.../services/org.apache.servicecomb.core.Transport | 1 +
51 files changed, 1073 insertions(+), 156 deletions(-)
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/CommonRestConfiguration.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/CommonRestConfiguration.java
index 5fb174b72..655c9b4f8 100644
---
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/CommonRestConfiguration.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/CommonRestConfiguration.java
@@ -26,6 +26,7 @@ import
org.apache.servicecomb.common.rest.codec.query.QueryCodecSsv;
import org.apache.servicecomb.common.rest.codec.query.QueryCodecs;
import org.apache.servicecomb.common.rest.codec.query.QueryCodecsUtils;
import org.apache.servicecomb.common.rest.filter.inner.RestServerCodecFilter;
+import
org.apache.servicecomb.common.rest.filter.inner.WebSocketServerCodecFilter;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@@ -61,6 +62,11 @@ public class CommonRestConfiguration {
return new RestServerCodecFilter();
}
+ @Bean
+ public WebSocketServerCodecFilter webSocketServerCodecFilter() {
+ return new WebSocketServerCodecFilter();
+ }
+
@Bean
public QueryCodecs queryCodecs(List<QueryCodec> orderedCodecs) {
return new QueryCodecs(orderedCodecs);
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/EdgeServerWebSocketInvocationCreator.java
similarity index 68%
copy from
edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
copy to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/EdgeServerWebSocketInvocationCreator.java
index 5379515af..ff0b0f381 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/EdgeServerWebSocketInvocationCreator.java
@@ -14,42 +14,33 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.edge.core;
+
+package org.apache.servicecomb.common.rest;
import java.util.concurrent.CompletableFuture;
-import org.apache.servicecomb.common.rest.RestVertxProducerInvocationCreator;
import org.apache.servicecomb.common.rest.locator.OperationLocator;
import org.apache.servicecomb.common.rest.locator.ServicePathManager;
-import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.Endpoint;
import org.apache.servicecomb.core.Invocation;
import org.apache.servicecomb.core.SCBEngine;
import org.apache.servicecomb.core.invocation.InvocationFactory;
import
org.apache.servicecomb.core.provider.consumer.MicroserviceReferenceConfig;
import org.apache.servicecomb.core.provider.consumer.ReferenceConfig;
-import org.apache.servicecomb.foundation.vertx.http.HttpServletRequestEx;
-import org.apache.servicecomb.foundation.vertx.http.HttpServletResponseEx;
-
-import io.vertx.core.Vertx;
-import io.vertx.ext.web.RoutingContext;
-public class EdgeInvocationCreator extends RestVertxProducerInvocationCreator {
- public static final String EDGE_INVOCATION_CONTEXT = "edgeInvocationContext";
+import io.vertx.core.http.HttpMethod;
+import io.vertx.core.http.ServerWebSocket;
- protected final String microserviceName;
-
- protected final String path;
+public class EdgeServerWebSocketInvocationCreator extends
ProviderServerWebSocketInvocationCreator {
+ private final String microserviceName;
protected MicroserviceReferenceConfig microserviceReferenceConfig;
- public EdgeInvocationCreator(RoutingContext routingContext,
- HttpServletRequestEx requestEx, HttpServletResponseEx responseEx,
- String microserviceName, String path) {
- // Set endpoint before load balance because edge service will use RESTFUL
transport filters.
- super(routingContext, null,
-
SCBEngine.getInstance().getTransportManager().findTransport(CoreConst.RESTFUL).getEndpoint(),
- requestEx, responseEx);
+ protected final String path;
+ public EdgeServerWebSocketInvocationCreator(String microserviceName, String
path,
+ Endpoint endpoint, ServerWebSocket webSocket) {
+ super(null, endpoint, webSocket);
this.microserviceName = microserviceName;
this.path = path;
}
@@ -71,7 +62,7 @@ public class EdgeInvocationCreator extends
RestVertxProducerInvocationCreator {
@Override
protected OperationLocator locateOperation(ServicePathManager
servicePathManager) {
- return servicePathManager.consumerLocateOperation(path,
requestEx.getMethod());
+ return servicePathManager.consumerLocateOperation(path,
HttpMethod.POST.name());
}
@Override
@@ -85,7 +76,7 @@ public class EdgeInvocationCreator extends
RestVertxProducerInvocationCreator {
null);
invocation.setSync(false);
invocation.setEdge();
- invocation.addLocalContext(EDGE_INVOCATION_CONTEXT,
Vertx.currentContext());
+ invocation.setEndpoint(endpoint); // ensure transport name is correct
return invocation;
}
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ProviderServerWebSocketInvocationCreator.java
similarity index 54%
copy from
core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
copy to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ProviderServerWebSocketInvocationCreator.java
index a413fce10..bcd4fc827 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ProviderServerWebSocketInvocationCreator.java
@@ -14,15 +14,24 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.core.invocation;
-import java.util.concurrent.CompletableFuture;
+package org.apache.servicecomb.common.rest;
+import org.apache.servicecomb.core.Endpoint;
import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.definition.MicroserviceMeta;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
-public interface InvocationCreator {
- CompletableFuture<Invocation> createAsync();
+import io.vertx.core.http.ServerWebSocket;
+
+public class ProviderServerWebSocketInvocationCreator extends
ServerWebSocketInvocationCreator {
+ public ProviderServerWebSocketInvocationCreator(MicroserviceMeta
microserviceMeta,
+ Endpoint endpoint, ServerWebSocket webSocket) {
+ super(microserviceMeta, endpoint, webSocket);
+ }
+
+ @Override
+ protected void initTransportContext(Invocation invocation) {
+ WebSocketTransportContext transportContext = new
WebSocketTransportContext(websocket);
+ invocation.setTransportContext(transportContext);
+ }
}
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ServerWebSocketInvocationCreator.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ServerWebSocketInvocationCreator.java
new file mode 100644
index 000000000..46d8b7051
--- /dev/null
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/ServerWebSocketInvocationCreator.java
@@ -0,0 +1,136 @@
+/*
+ * 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.servicecomb.common.rest;
+
+import static jakarta.ws.rs.core.Response.Status.NOT_FOUND;
+import static
org.apache.servicecomb.core.exception.ExceptionCodes.NOT_DEFINED_ANY_SCHEMA;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+
+import org.apache.commons.lang3.StringUtils;
+import org.apache.servicecomb.common.rest.definition.RestOperationMeta;
+import org.apache.servicecomb.common.rest.locator.OperationLocator;
+import org.apache.servicecomb.common.rest.locator.ServicePathManager;
+import org.apache.servicecomb.config.YAMLUtil;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.Endpoint;
+import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.definition.MicroserviceMeta;
+import org.apache.servicecomb.core.exception.Exceptions;
+import org.apache.servicecomb.core.invocation.InvocationCreator;
+import org.apache.servicecomb.core.invocation.InvocationFactory;
+import org.apache.servicecomb.foundation.common.LegacyPropertyFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.vertx.core.http.HttpMethod;
+import io.vertx.core.http.ServerWebSocket;
+import io.vertx.core.json.Json;
+
+public abstract class ServerWebSocketInvocationCreator implements
InvocationCreator {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(RestVertxProducerInvocationCreator.class);
+
+ protected MicroserviceMeta microserviceMeta;
+
+ protected final Endpoint endpoint;
+
+ protected final ServerWebSocket websocket;
+
+ protected RestOperationMeta restOperationMeta;
+
+ public ServerWebSocketInvocationCreator(MicroserviceMeta microserviceMeta,
Endpoint endpoint,
+ ServerWebSocket websocket) {
+ this.microserviceMeta = microserviceMeta;
+ this.endpoint = endpoint;
+ this.websocket = websocket;
+ }
+
+ @Override
+ public CompletableFuture<Invocation> createAsync() {
+ initRestOperation();
+
+ Invocation invocation = createInstance();
+ initInvocationContext(invocation);
+ addParameterContext(invocation);
+ initTransportContext(invocation);
+
+ return CompletableFuture.completedFuture(invocation);
+ }
+
+ protected Invocation createInstance() {
+ return InvocationFactory.forProvider(endpoint,
restOperationMeta.getOperationMeta(), null);
+ }
+
+ protected void initInvocationContext(Invocation invocation) {
+ if
(!LegacyPropertyFactory.getBooleanProperty(RestConst.DECODE_INVOCATION_CONTEXT,
true)) {
+ return;
+ }
+
+ String strCseContext = websocket.headers().get(CoreConst.CSE_CONTEXT);
+ if (StringUtils.isEmpty(strCseContext)) {
+ return;
+ }
+
+ @SuppressWarnings("unchecked")
+ Map<String, String> invocationContext = Json.decodeValue(strCseContext,
Map.class);
+ invocation.mergeContext(invocationContext);
+ }
+
+ // No queries for websocket
+ protected void addParameterContext(Invocation invocation) {
+ String headerContextMapper = LegacyPropertyFactory
+ .getStringProperty(RestConst.HEADER_CONTEXT_MAPPER);
+
+ Map<String, Object> headerContextMappers;
+ if (headerContextMapper != null) {
+ headerContextMappers = YAMLUtil.yaml2Properties(headerContextMapper);
+ } else {
+ headerContextMappers = new HashMap<>();
+ }
+
+ headerContextMappers.forEach((k, v) -> {
+ if (v instanceof String && websocket.headers().get(k) != null) {
+ invocation.addContext((String) v, websocket.headers().get(k));
+ }
+ });
+ }
+
+ protected abstract void initTransportContext(Invocation invocation);
+
+ protected void initRestOperation() {
+ OperationLocator locator = locateOperation(microserviceMeta);
+ restOperationMeta = locator.getOperation();
+ }
+
+ protected OperationLocator locateOperation(MicroserviceMeta
microserviceMeta) {
+ ServicePathManager servicePathManager =
ServicePathManager.getServicePathManager(microserviceMeta);
+ if (servicePathManager == null) {
+ LOGGER.error("No schema defined for {}:{}.",
this.microserviceMeta.getAppId(),
+ this.microserviceMeta.getMicroserviceName());
+ throw Exceptions.create(NOT_FOUND, NOT_DEFINED_ANY_SCHEMA,
NOT_FOUND.getReasonPhrase());
+ }
+
+ return locateOperation(servicePathManager);
+ }
+
+ protected OperationLocator locateOperation(ServicePathManager
servicePathManager) {
+ return servicePathManager.producerLocateOperation(websocket.path(),
HttpMethod.POST.name());
+ }
+}
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/WebSocketTransportContext.java
similarity index 62%
copy from
core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
copy to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/WebSocketTransportContext.java
index a413fce10..738b09788 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/WebSocketTransportContext.java
@@ -14,15 +14,21 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.core.invocation;
-import java.util.concurrent.CompletableFuture;
+package org.apache.servicecomb.common.rest;
-import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.swagger.invocation.context.TransportContext;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
-public interface InvocationCreator {
- CompletableFuture<Invocation> createAsync();
+import io.vertx.core.http.ServerWebSocket;
+
+public class WebSocketTransportContext implements TransportContext {
+ private final ServerWebSocket serverWebSocket;
+
+ public WebSocketTransportContext(ServerWebSocket serverWebSocket) {
+ this.serverWebSocket = serverWebSocket;
+ }
+
+ public ServerWebSocket getServerWebSocket() {
+ return this.serverWebSocket;
+ }
}
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/WebSocketServerCodecFilter.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/WebSocketServerCodecFilter.java
new file mode 100644
index 000000000..79e736852
--- /dev/null
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/WebSocketServerCodecFilter.java
@@ -0,0 +1,124 @@
+/*
+ * 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.servicecomb.common.rest.filter.inner;
+
+import static
org.apache.servicecomb.core.exception.Exceptions.toProducerResponse;
+
+import java.util.HashMap;
+import java.util.concurrent.CompletableFuture;
+
+import org.apache.servicecomb.common.rest.WebSocketTransportContext;
+import org.apache.servicecomb.common.rest.codec.produce.ProduceJsonProcessor;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.filter.AbstractFilter;
+import org.apache.servicecomb.core.filter.EdgeFilter;
+import org.apache.servicecomb.core.filter.Filter;
+import org.apache.servicecomb.core.filter.FilterNode;
+import org.apache.servicecomb.core.filter.ProviderFilter;
+import org.apache.servicecomb.foundation.vertx.stream.BufferOutputStream;
+import org.apache.servicecomb.swagger.invocation.Response;
+import org.apache.servicecomb.swagger.invocation.context.TransportContext;
+import org.apache.servicecomb.swagger.invocation.exception.InvocationException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.vertx.core.buffer.Buffer;
+import io.vertx.core.http.ServerWebSocket;
+
+public class WebSocketServerCodecFilter extends AbstractFilter implements
ProviderFilter, EdgeFilter {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(WebSocketServerCodecFilter.class);
+
+ public static final String NAME = "websocket-codec";
+
+ @Override
+ public String getName() {
+ return NAME;
+ }
+
+ @Override
+ public int getOrder() {
+ // almost time, should be the first filter.
+ return Filter.PROVIDER_SCHEDULE_FILTER_ORDER - 2000;
+ }
+
+ @Override
+ public boolean enabledForTransport(String transport) {
+ return CoreConst.WEBSOCKET.equals(transport);
+ }
+
+ @Override
+ public CompletableFuture<Response> onFilter(Invocation invocation,
FilterNode nextNode) {
+ return CompletableFuture.completedFuture(invocation)
+ .thenCompose(this::decodeRequest)
+ .thenCompose(v -> invokeNext(invocation, nextNode))
+ .exceptionally(exception -> toProducerResponse(invocation, exception))
+ .thenCompose(response -> encodeResponse(invocation, response));
+ }
+
+ protected CompletableFuture<Response> invokeNext(Invocation invocation,
FilterNode nextNode) {
+ if (invocation.isEdge()) {
+ TransportContext transportContext = invocation.getTransportContext();
+ return nextNode.onFilter(invocation).whenComplete((r, e) ->
invocation.setTransportContext(transportContext));
+ }
+ return nextNode.onFilter(invocation);
+ }
+
+ protected CompletableFuture<Void> decodeRequest(Invocation invocation) {
+ invocation.getInvocationStageTrace().startProviderDecodeRequest();
+ invocation.setSwaggerArguments(new HashMap<>()); // set context parameters
and do nothing else.
+ invocation.getInvocationStageTrace().finishProviderDecodeRequest();
+ return CompletableFuture.completedFuture(null);
+ }
+
+ protected CompletableFuture<Response> encodeResponse(Invocation invocation,
Response response) {
+ invocation.onEncodeResponseStart(response);
+ WebSocketTransportContext context = invocation.getTransportContext();
+
+ return encodeResponse(response, context.getServerWebSocket())
+ .whenComplete((r, e) -> invocation.onEncodeResponseFinish());
+ }
+
+ private static boolean isFailedResponse(Response response) {
+ return response.getResult() instanceof InvocationException;
+ }
+
+ private static CompletableFuture<Response> writeResponse(
+ ServerWebSocket webSocket, Object data, Response response) {
+ try (BufferOutputStream output = new BufferOutputStream(Buffer.buffer())) {
+ ProduceJsonProcessor produceProcessor = new ProduceJsonProcessor();
+ produceProcessor.encodeResponse(output, data);
+ CompletableFuture<Response> result = new CompletableFuture<>();
+ webSocket.write(output.getBuffer()).onComplete(v ->
+ result.complete(response), result::completeExceptionally);
+
+ return result;
+ } catch (Throwable e) {
+ LOGGER.error("internal service error must be fixed.", e);
+ return CompletableFuture.failedFuture(e);
+ }
+ }
+
+ public static CompletableFuture<Response> encodeResponse(Response response,
ServerWebSocket webSocket) {
+ if (isFailedResponse(response)) {
+ return writeResponse(webSocket, ((InvocationException)
response.getResult()).getErrorData(),
+ response);
+ }
+ return CompletableFuture.completedFuture(response);
+ }
+}
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationItem.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationItem.java
similarity index 97%
rename from
edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationItem.java
rename to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationItem.java
index 2fa38461f..d884e49e4 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationItem.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationItem.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.servicecomb.edge.core;
+package org.apache.servicecomb.common.rest.route;
import java.util.regex.Pattern;
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationLoader.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationLoader.java
similarity index 98%
rename from
edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationLoader.java
rename to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationLoader.java
index ae93c1349..5c3df2c33 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedConfigurationLoader.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/URLMappedConfigurationLoader.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.servicecomb.edge.core;
+package org.apache.servicecomb.common.rest.route;
import java.util.HashMap;
import java.util.Map;
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/Utils.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/Utils.java
similarity index 96%
rename from
edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/Utils.java
rename to
common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/Utils.java
index 9bbefb9f0..0e1f64291 100644
--- a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/Utils.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/route/Utils.java
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-package org.apache.servicecomb.edge.core;
+package org.apache.servicecomb.common.rest.route;
/**
* Commonly used methods in this package.
diff --git a/core/src/main/java/org/apache/servicecomb/core/CoreConst.java
b/core/src/main/java/org/apache/servicecomb/core/CoreConst.java
index dfe2aeebf..37c34b9c9 100644
--- a/core/src/main/java/org/apache/servicecomb/core/CoreConst.java
+++ b/core/src/main/java/org/apache/servicecomb/core/CoreConst.java
@@ -31,6 +31,8 @@ public final class CoreConst {
public static final String HIGHWAY = "highway";
+ public static final String WEBSOCKET = "websocket";
+
public static final String ANY_TRANSPORT = "";
public static final String VERSION_RULE_LATEST =
DefinitionConst.VERSION_RULE_LATEST;
diff --git a/core/src/main/java/org/apache/servicecomb/core/Transport.java
b/core/src/main/java/org/apache/servicecomb/core/Transport.java
index 4777a5ca0..1c40303fc 100644
--- a/core/src/main/java/org/apache/servicecomb/core/Transport.java
+++ b/core/src/main/java/org/apache/servicecomb/core/Transport.java
@@ -19,7 +19,6 @@ package org.apache.servicecomb.core;
import org.springframework.core.env.Environment;
-// TODO:感觉要拆成显式的client、server才好些
public interface Transport {
String getName();
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
index a413fce10..323ad924a 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
@@ -20,9 +20,6 @@ import java.util.concurrent.CompletableFuture;
import org.apache.servicecomb.core.Invocation;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
public interface InvocationCreator {
CompletableFuture<Invocation> createAsync();
}
diff --git
a/demo/demo-zookeeper/consumer/src/main/java/org/apache/servicecomb/samples/ClientWebsocketController.java
b/demo/demo-zookeeper/consumer/src/main/java/org/apache/servicecomb/samples/ClientWebsocketController.java
new file mode 100644
index 000000000..76fb2fb17
--- /dev/null
+++
b/demo/demo-zookeeper/consumer/src/main/java/org/apache/servicecomb/samples/ClientWebsocketController.java
@@ -0,0 +1,54 @@
+/*
+ * 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.servicecomb.samples;
+
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.annotation.Transport;
+import org.apache.servicecomb.provider.pojo.RpcReference;
+import org.apache.servicecomb.provider.rest.common.RestSchema;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+
+import io.vertx.core.http.ServerWebSocket;
+import io.vertx.core.http.WebSocket;
+
+@RestSchema(schemaId = "ClientWebsocketController")
+@RequestMapping(path = "/ws")
+public class ClientWebsocketController {
+ interface ProviderService {
+ WebSocket websocket();
+ }
+
+ @RpcReference(schemaId = "WebsocketController", microserviceName =
"provider")
+ private ProviderService providerService;
+
+ @PostMapping("/websocket")
+ @Transport(name = CoreConst.WEBSOCKET)
+ public void websocket(ServerWebSocket serverWebsocket) {
+ WebSocket providerWebSocket = providerService.websocket();
+ providerWebSocket.closeHandler(v -> serverWebsocket.close());
+ providerWebSocket.textMessageHandler(m -> {
+ System.out.println("send message " + m);
+ serverWebsocket.writeTextMessage(m);
+ });
+ serverWebsocket.textMessageHandler(m -> {
+ System.out.println("receive message " + m);
+ providerWebSocket.writeTextMessage(m);
+ });
+ }
+}
diff --git a/demo/demo-zookeeper/consumer/src/main/resources/application.yml
b/demo/demo-zookeeper/consumer/src/main/resources/application.yml
index a0e95d3ec..e3f107bd8 100644
--- a/demo/demo-zookeeper/consumer/src/main/resources/application.yml
+++ b/demo/demo-zookeeper/consumer/src/main/resources/application.yml
@@ -28,7 +28,9 @@ servicecomb:
connectString: 127.0.0.1:2181
rest:
- address: 0.0.0.0:9092
+ address: 0.0.0.0:9092?websocketEnabled=true
+ server:
+ websocket-prefix: /ws
highway:
address: 0.0.0.0:7092
diff --git a/demo/demo-zookeeper/gateway/src/main/resources/application.yml
b/demo/demo-zookeeper/gateway/src/main/resources/application.yml
index f63d61c27..ea13ef73d 100644
--- a/demo/demo-zookeeper/gateway/src/main/resources/application.yml
+++ b/demo/demo-zookeeper/gateway/src/main/resources/application.yml
@@ -26,7 +26,9 @@ servicecomb:
connectString: 127.0.0.1:2181
rest:
- address: 0.0.0.0:9090?sslEnabled=false
+ address: 0.0.0.0:9090?websocketEnabled=true
+ server:
+ websocket-prefix: /ws
http:
dispatcher:
@@ -42,3 +44,10 @@ servicecomb:
path: "/.*"
microserviceName: consumer
versionRule: 0.0.0+
+ websocket:
+ mappings:
+ consumer:
+ prefixSegmentCount: 0
+ path: "/ws/.*"
+ microserviceName: consumer
+ versionRule: 0.0.0+
diff --git
a/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/WebsocketController.java
b/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/WebsocketController.java
new file mode 100644
index 000000000..9ecf678a1
--- /dev/null
+++
b/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/WebsocketController.java
@@ -0,0 +1,57 @@
+/*
+ * 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.servicecomb.samples;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.annotation.Transport;
+import org.apache.servicecomb.provider.rest.common.RestSchema;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+
+import io.vertx.core.http.ServerWebSocket;
+
+@RestSchema(schemaId = "WebsocketController")
+@RequestMapping(path = "/ws")
+public class WebsocketController {
+ @PostMapping("/websocket")
+ @Transport(name = CoreConst.WEBSOCKET)
+ public void websocket(ServerWebSocket serverWebsocket) {
+ AtomicInteger receiveCount = new AtomicInteger(0);
+ serverWebsocket.writeTextMessage("hello", r -> {
+ });
+ serverWebsocket.textMessageHandler(s -> {
+ receiveCount.getAndIncrement();
+ });
+ serverWebsocket.closeHandler((v) -> System.out.println("closed"));
+ new Thread(() -> {
+ for (int i = 0; i < 5; i++) {
+ serverWebsocket.writeTextMessage("hello " + i, r -> {
+ });
+ try {
+ Thread.sleep(500);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ }
+ serverWebsocket.writeTextMessage("total " + receiveCount.get());
+ serverWebsocket.close();
+ }).start();
+ }
+}
diff --git a/demo/demo-zookeeper/provider/src/main/resources/application.yml
b/demo/demo-zookeeper/provider/src/main/resources/application.yml
index 10e0ab805..dd777236b 100644
--- a/demo/demo-zookeeper/provider/src/main/resources/application.yml
+++ b/demo/demo-zookeeper/provider/src/main/resources/application.yml
@@ -32,7 +32,9 @@ servicecomb:
instance-tag: config-demo
rest:
- address: 0.0.0.0:9094
+ address: 0.0.0.0:9094?websocketEnabled=true
+ server:
+ websocket-prefix: /ws
highway:
address: 0.0.0.0:7094
diff --git
a/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/ThirdSvcConfiguration.java
b/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/ThirdSvcConfiguration.java
index ddb65d14a..5f23c66cd 100644
---
a/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/ThirdSvcConfiguration.java
+++
b/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/ThirdSvcConfiguration.java
@@ -19,6 +19,8 @@ package org.apache.servicecomb.samples;
import java.util.List;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.annotation.Transport;
import org.apache.servicecomb.localregistry.RegistryBean;
import org.apache.servicecomb.localregistry.RegistryBean.Instance;
import org.apache.servicecomb.localregistry.RegistryBean.Instances;
@@ -27,10 +29,20 @@ import org.reactivestreams.Publisher;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
+import io.vertx.core.http.WebSocket;
+
@Configuration
public class ThirdSvcConfiguration {
+ @RequestMapping(path = "/ws")
+ public interface WebsocketClient {
+ @PostMapping("/websocket")
+ @Transport(name = CoreConst.WEBSOCKET)
+ WebSocket websocket();
+ }
+
@RequestMapping(path = "/")
public interface ReactiveStreamClient {
class Model {
@@ -88,11 +100,12 @@ public class ThirdSvcConfiguration {
public RegistryBean gatewayServiceBean() {
return new RegistryBean()
.addSchemaInterface("ReactiveStreamController",
ReactiveStreamClient.class)
+ .addSchemaInterface("WebsocketController", WebsocketClient.class)
.setAppId("demo-zookeeper")
.setServiceName("gateway")
.setVersion("0.0.1")
.setInstances(new Instances().setInstances(List.of(
- new Instance().setEndpoints(List.of("rest://localhost:9090")))));
+ new
Instance().setEndpoints(List.of("rest://localhost:9090?websocketEnabled=true")))));
}
@Bean("reactiveStreamProvider")
@@ -104,4 +117,9 @@ public class ThirdSvcConfiguration {
public ReactiveStreamClient reactiveStreamGateway() {
return Invoker.createProxy("gateway", "ReactiveStreamController",
ReactiveStreamClient.class);
}
+
+ @Bean
+ public WebsocketClient gatewayWebsocketClient() {
+ return Invoker.createProxy("gateway", "WebsocketController",
WebsocketClient.class);
+ }
}
diff --git
a/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/WebsocketIT.java
b/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/WebsocketIT.java
new file mode 100644
index 000000000..ea7a7f45b
--- /dev/null
+++
b/demo/demo-zookeeper/test-client/src/main/java/org/apache/servicecomb/samples/WebsocketIT.java
@@ -0,0 +1,57 @@
+/*
+ * 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.servicecomb.samples;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.servicecomb.demo.CategorizedTestCase;
+import org.apache.servicecomb.demo.TestMgr;
+import org.apache.servicecomb.samples.ThirdSvcConfiguration.WebsocketClient;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import io.vertx.core.http.WebSocket;
+
+@Component
+public class WebsocketIT implements CategorizedTestCase {
+ @Autowired
+ private WebsocketClient websocketClient;
+
+ @Override
+ public void testRestTransport() throws Exception {
+ StringBuffer sb = new StringBuffer();
+ AtomicBoolean closed = new AtomicBoolean(false);
+ CountDownLatch latch = new CountDownLatch(1);
+
+ WebSocket webSocket = websocketClient.websocket();
+ webSocket.textMessageHandler(s -> {
+ sb.append(s);
+ sb.append(" ");
+ webSocket.writeTextMessage(s);
+ });
+ webSocket.closeHandler(v -> {
+ closed.set(true);
+ latch.countDown();
+ });
+ latch.await(30, TimeUnit.SECONDS);
+ TestMgr.check(sb.toString(), "hello hello 0 hello 1 hello 2 hello 3 hello
4 total 6 ");
+ TestMgr.check(closed.get(), true);
+ }
+}
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/CommonHttpEdgeDispatcher.java
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/CommonHttpEdgeDispatcher.java
index d771fa80a..0f7e59bb0 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/CommonHttpEdgeDispatcher.java
+++
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/CommonHttpEdgeDispatcher.java
@@ -20,6 +20,9 @@ package org.apache.servicecomb.edge.core;
import java.util.HashMap;
import java.util.Map;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationItem;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationLoader;
+import org.apache.servicecomb.common.rest.route.Utils;
import org.apache.servicecomb.config.BootStrapProperties;
import org.apache.servicecomb.config.ConfigurationChangedEvent;
import org.apache.servicecomb.core.Invocation;
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/DefaultEdgeDispatcher.java
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/DefaultEdgeDispatcher.java
index 08c33b2a8..56108a8cf 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/DefaultEdgeDispatcher.java
+++
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/DefaultEdgeDispatcher.java
@@ -18,6 +18,7 @@
package org.apache.servicecomb.edge.core;
import org.apache.servicecomb.common.rest.RestProducerInvocationFlow;
+import org.apache.servicecomb.common.rest.route.Utils;
import org.apache.servicecomb.core.invocation.InvocationCreator;
import org.apache.servicecomb.foundation.common.LegacyPropertyFactory;
import org.apache.servicecomb.foundation.vertx.http.HttpServletRequestEx;
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeBootListener.java
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeBootListener.java
index 538c83b8b..144dc2c62 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeBootListener.java
+++
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeBootListener.java
@@ -20,21 +20,8 @@ package org.apache.servicecomb.edge.core;
import org.apache.servicecomb.core.BootListener;
import org.apache.servicecomb.core.executor.ExecutorManager;
import org.apache.servicecomb.transport.rest.vertx.TransportConfig;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.core.env.Environment;
public class EdgeBootListener implements BootListener {
- private static final Logger LOGGER =
LoggerFactory.getLogger(EdgeBootListener.class);
-
- private Environment environment;
-
- @Autowired
- public void setEnvironment(Environment environment) {
- this.environment = environment;
- }
-
@Override
public void onBootEvent(BootEvent event) {
if (!EventType.BEFORE_PRODUCER_PROVIDER.equals(event.getEventType())) {
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
index 5379515af..309548a5b 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
+++
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/EdgeInvocationCreator.java
@@ -85,6 +85,7 @@ public class EdgeInvocationCreator extends
RestVertxProducerInvocationCreator {
null);
invocation.setSync(false);
invocation.setEdge();
+ invocation.setEndpoint(endpoint); // ensure transport name is correct
invocation.addLocalContext(EDGE_INVOCATION_CONTEXT,
Vertx.currentContext());
return invocation;
diff --git
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedEdgeDispatcher.java
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedEdgeDispatcher.java
index 195513714..d4e456c09 100644
---
a/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedEdgeDispatcher.java
+++
b/edge/edge-core/src/main/java/org/apache/servicecomb/edge/core/URLMappedEdgeDispatcher.java
@@ -21,16 +21,18 @@ import java.util.HashMap;
import java.util.Map;
import org.apache.servicecomb.common.rest.RestProducerInvocationFlow;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationItem;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationLoader;
+import org.apache.servicecomb.common.rest.route.Utils;
import org.apache.servicecomb.config.ConfigurationChangedEvent;
import org.apache.servicecomb.core.invocation.InvocationCreator;
import org.apache.servicecomb.foundation.common.LegacyPropertyFactory;
+import org.apache.servicecomb.foundation.common.event.EventManager;
import org.apache.servicecomb.foundation.vertx.http.HttpServletRequestEx;
import org.apache.servicecomb.foundation.vertx.http.HttpServletResponseEx;
import
org.apache.servicecomb.foundation.vertx.http.VertxServerRequestToHttpServletRequest;
import
org.apache.servicecomb.foundation.vertx.http.VertxServerResponseToHttpServletResponse;
import org.apache.servicecomb.transport.rest.vertx.RestBodyHandler;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.core.env.Environment;
@@ -45,8 +47,6 @@ import io.vertx.ext.web.handler.PlatformHandler;
* Provide a URL mapping based dispatcher. Users configure witch URL patterns
dispatch to a target service.
*/
public class URLMappedEdgeDispatcher extends AbstractEdgeDispatcher {
- private static final Logger LOG =
LoggerFactory.getLogger(URLMappedEdgeDispatcher.class);
-
public static final String CONFIGURATION_ITEM = "URLMappedConfigurationItem";
private static final String PATTERN_ANY = "/(.*)";
@@ -64,6 +64,7 @@ public class URLMappedEdgeDispatcher extends
AbstractEdgeDispatcher {
private Environment environment;
public URLMappedEdgeDispatcher() {
+ EventManager.register(this);
}
// though this is an SPI, but add as beans.
diff --git
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestURLMappedEdgeDispatcher.java
b/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestURLMappedEdgeDispatcher.java
index 2e23ab673..046adfaa3 100644
---
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestURLMappedEdgeDispatcher.java
+++
b/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestURLMappedEdgeDispatcher.java
@@ -20,6 +20,7 @@ package org.apache.servicecomb.edge.core;
import java.util.HashMap;
import java.util.Map;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationItem;
import org.apache.servicecomb.config.ConfigurationChangedEvent;
import org.apache.servicecomb.transport.rest.vertx.RestBodyHandler;
import org.junit.jupiter.api.AfterEach;
diff --git
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
b/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
index f7fcfa18c..38e8fa837 100644
---
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
+++
b/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
@@ -17,6 +17,7 @@
package org.apache.servicecomb.edge.core;
+import org.apache.servicecomb.common.rest.route.Utils;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
diff --git
a/foundations/foundation-common/src/main/java/org/apache/servicecomb/foundation/common/net/URIEndpointObject.java
b/foundations/foundation-common/src/main/java/org/apache/servicecomb/foundation/common/net/URIEndpointObject.java
index 03b0ffc96..7d4b8a00b 100644
---
a/foundations/foundation-common/src/main/java/org/apache/servicecomb/foundation/common/net/URIEndpointObject.java
+++
b/foundations/foundation-common/src/main/java/org/apache/servicecomb/foundation/common/net/URIEndpointObject.java
@@ -35,16 +35,23 @@ public class URIEndpointObject extends IpPort {
private static final String PROTOCOL_KEY = "protocol";
+ private static final String WEBSOCKET_ENABLED_KEY = "websocketEnabled";
+
private static final String HTTP2 = "http2";
private final boolean sslEnabled;
private boolean http2Enabled;
+ private boolean websocketEnabled;
+
private final Map<String, List<String>> querys;
+ private final String schema;
+
public URIEndpointObject(String endpoint) {
URI uri = URI.create(endpoint);
+ schema = uri.getScheme();
setHostOrIp(uri.getHost());
if (uri.getPort() < 0) {
// do not use default port
@@ -53,6 +60,7 @@ public class URIEndpointObject extends IpPort {
setPort(uri.getPort());
querys = splitQuery(uri);
sslEnabled = Boolean.parseBoolean(getFirst(SSL_ENABLED_KEY));
+ websocketEnabled = Boolean.parseBoolean(getFirst(WEBSOCKET_ENABLED_KEY));
String httpVersion = getFirst(PROTOCOL_KEY);
if (HTTP2.equals(httpVersion)) {
http2Enabled = true;
@@ -73,10 +81,18 @@ public class URIEndpointObject extends IpPort {
return sslEnabled;
}
+ public boolean isWebsocketEnabled() {
+ return websocketEnabled;
+ }
+
public boolean isHttp2Enabled() {
return http2Enabled;
}
+ public String getSchema() {
+ return this.schema;
+ }
+
public List<String> getQuery(String key) {
return querys.get(key);
}
diff --git
a/foundations/foundation-registry/src/main/java/org/apache/servicecomb/registry/discovery/AbstractEndpointDiscoveryFilter.java
b/foundations/foundation-registry/src/main/java/org/apache/servicecomb/registry/discovery/AbstractEndpointDiscoveryFilter.java
index 45a20e74a..f4b464f00 100644
---
a/foundations/foundation-registry/src/main/java/org/apache/servicecomb/registry/discovery/AbstractEndpointDiscoveryFilter.java
+++
b/foundations/foundation-registry/src/main/java/org/apache/servicecomb/registry/discovery/AbstractEndpointDiscoveryFilter.java
@@ -17,10 +17,10 @@
package org.apache.servicecomb.registry.discovery;
-import java.net.URI;
import java.util.ArrayList;
import java.util.List;
+import org.apache.servicecomb.foundation.common.net.URIEndpointObject;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -29,6 +29,10 @@ public abstract class AbstractEndpointDiscoveryFilter
implements DiscoveryFilter
private static final String ALL_TRANSPORT = "";
+ private static final String WEBSOCKET_TRANSPORT = "websocket";
+
+ private static final String REST_TRANSPORT = "rest";
+
@Override
public boolean isGroupingFilter() {
return true;
@@ -48,13 +52,16 @@ public abstract class AbstractEndpointDiscoveryFilter
implements DiscoveryFilter
for (StatefulDiscoveryInstance instance : instances) {
for (String endpoint : instance.getEndpoints()) {
try {
- URI uri = URI.create(endpoint);
- String transportName = uri.getScheme();
- if (!isTransportNameMatch(transportName, expectTransportName)) {
+ URIEndpointObject endpointObject = new URIEndpointObject(endpoint);
+ String transPortName = endpointObject.getSchema();
+ if (endpointObject.isWebsocketEnabled() &&
WEBSOCKET_TRANSPORT.equals(expectTransportName)) {
+ transPortName = WEBSOCKET_TRANSPORT;
+ }
+ if (!isTransportNameMatch(transPortName, expectTransportName)) {
continue;
}
- Object objEndpoint = createEndpoint(context, transportName,
endpoint, instance);
+ Object objEndpoint = createEndpoint(context, transPortName,
endpoint, instance);
if (objEndpoint == null) {
continue;
}
@@ -72,7 +79,10 @@ public abstract class AbstractEndpointDiscoveryFilter
implements DiscoveryFilter
}
protected boolean isTransportNameMatch(String transportName, String
expectTransportName) {
- return ALL_TRANSPORT.equals(expectTransportName) ||
transportName.equals(expectTransportName);
+ if (ALL_TRANSPORT.equals(expectTransportName)) {
+ return true;
+ }
+ return transportName.equals(expectTransportName);
}
protected abstract String findTransportName(DiscoveryContext context,
DiscoveryTreeNode parent);
diff --git
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/VertxTLSBuilder.java
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/VertxTLSBuilder.java
index e83a6fec9..ba2203338 100644
---
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/VertxTLSBuilder.java
+++
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/VertxTLSBuilder.java
@@ -33,6 +33,7 @@ import org.slf4j.LoggerFactory;
import io.vertx.core.http.ClientAuth;
import io.vertx.core.http.HttpClientOptions;
+import io.vertx.core.http.WebSocketClientOptions;
import io.vertx.core.net.ClientOptionsBase;
import io.vertx.core.net.JksOptions;
import io.vertx.core.net.NetServerOptions;
@@ -87,6 +88,18 @@ public final class VertxTLSBuilder {
buildHttpClientOptions(sslOption, sslCustom, httpClientOptions);
}
+ public static void buildWebSocketClientOptions(String sslKey,
WebSocketClientOptions webSocketClientOptions) {
+ SSLOptionFactory factory = SSLOptionFactory.createSSLOptionFactory(sslKey,
LegacyPropertyFactory.getEnvironment());
+ SSLOption sslOption;
+ if (factory == null) {
+ sslOption = SSLOption.build(sslKey,
LegacyPropertyFactory.getEnvironment());
+ } else {
+ sslOption = factory.createSSLOption();
+ }
+ SSLCustom sslCustom =
SSLCustom.createSSLCustom(sslOption.getSslCustomClass());
+ buildWebSocketClientOptions(sslOption, sslCustom, webSocketClientOptions);
+ }
+
public static HttpClientOptions buildHttpClientOptions(SSLOption sslOption,
SSLCustom sslCustom,
HttpClientOptions httpClientOptions) {
buildClientOptionsBase(sslOption, sslCustom, httpClientOptions);
@@ -94,6 +107,13 @@ public final class VertxTLSBuilder {
return httpClientOptions;
}
+ public static WebSocketClientOptions buildWebSocketClientOptions(SSLOption
sslOption, SSLCustom sslCustom,
+ WebSocketClientOptions webSocketClientOptions) {
+ buildClientOptionsBase(sslOption, sslCustom, webSocketClientOptions);
+ webSocketClientOptions.setVerifyHost(sslOption.isCheckCNHost());
+ return webSocketClientOptions;
+ }
+
public static ClientOptionsBase buildClientOptionsBase(SSLOption sslOption,
SSLCustom sslCustom,
ClientOptionsBase clientOptionsBase) {
buildTCPSSLOptions(sslOption, sslCustom, clientOptionsBase);
diff --git
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClientOptionsSPI.java
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClientOptionsSPI.java
index 78a12d099..8094cfeca 100644
---
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClientOptionsSPI.java
+++
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClientOptionsSPI.java
@@ -22,6 +22,8 @@ import
org.apache.servicecomb.foundation.vertx.VertxTLSBuilder;
import io.vertx.core.http.HttpClientOptions;
import io.vertx.core.http.HttpVersion;
+import io.vertx.core.http.WebSocketClientOptions;
+import io.vertx.core.net.ClientOptionsBase;
import io.vertx.core.net.ProxyOptions;
/**
@@ -97,17 +99,9 @@ public interface HttpClientOptionsSPI {
/***************** ssl settings ***************************/
boolean isSsl();
- static HttpClientOptions createHttpClientOptions(HttpClientOptionsSPI spi) {
- HttpClientOptions httpClientOptions = new HttpClientOptions();
-
- httpClientOptions.setProtocolVersion(spi.getHttpVersion());
+ static void buildClientOptionsBase(HttpClientOptionsSPI spi,
ClientOptionsBase httpClientOptions) {
httpClientOptions.setConnectTimeout(spi.getConnectTimeoutInMillis());
httpClientOptions.setIdleTimeout(spi.getIdleTimeoutInSeconds());
- httpClientOptions.setTryUseCompression(spi.isTryUseCompression());
- httpClientOptions.setMaxWaitQueueSize(spi.getMaxWaitQueueSize());
- httpClientOptions.setMaxPoolSize(spi.getMaxPoolSize());
- httpClientOptions.setKeepAlive(spi.isKeepAlive());
- httpClientOptions.setMaxHeaderSize(spi.getMaxHeaderSize());
httpClientOptions.setLogActivity(spi.enableLogActivity());
if (spi.isProxyEnable()) {
@@ -122,6 +116,21 @@ public interface HttpClientOptionsSPI {
if (spi.getHttpVersion() == HttpVersion.HTTP_2) {
httpClientOptions.setUseAlpn(spi.isUseAlpn());
+ }
+ }
+
+ static HttpClientOptions createHttpClientOptions(HttpClientOptionsSPI spi) {
+ HttpClientOptions httpClientOptions = new HttpClientOptions();
+ buildClientOptionsBase(spi, httpClientOptions);
+
+ httpClientOptions.setProtocolVersion(spi.getHttpVersion());
+ httpClientOptions.setTryUseCompression(spi.isTryUseCompression());
+ httpClientOptions.setMaxWaitQueueSize(spi.getMaxWaitQueueSize());
+ httpClientOptions.setMaxPoolSize(spi.getMaxPoolSize());
+ httpClientOptions.setKeepAlive(spi.isKeepAlive());
+ httpClientOptions.setMaxHeaderSize(spi.getMaxHeaderSize());
+
+ if (spi.getHttpVersion() == HttpVersion.HTTP_2) {
httpClientOptions.setHttp2ClearTextUpgrade(false);
httpClientOptions.setHttp2MultiplexingLimit(spi.getHttp2MultiplexingLimit());
httpClientOptions.setHttp2MaxPoolSize(spi.getHttp2MaxPoolSize());
@@ -136,4 +145,15 @@ public interface HttpClientOptionsSPI {
return httpClientOptions;
}
+
+ static WebSocketClientOptions
createWebSocketClientOptions(HttpClientOptionsSPI spi, boolean sslEnabled) {
+ WebSocketClientOptions webSocketClientOptions = new
WebSocketClientOptions();
+ buildClientOptionsBase(spi, webSocketClientOptions);
+
+ if (sslEnabled) {
+ VertxTLSBuilder.buildWebSocketClientOptions(spi.getConfigTag(),
webSocketClientOptions);
+ }
+
+ return webSocketClientOptions;
+ }
}
diff --git
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClients.java
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClients.java
index c3c898c8d..39d78422d 100644
---
a/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClients.java
+++
b/foundations/foundation-vertx/src/main/java/org/apache/servicecomb/foundation/vertx/client/http/HttpClients.java
@@ -36,6 +36,7 @@ import io.vertx.core.DeploymentOptions;
import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.core.dns.AddressResolverOptions;
+import io.vertx.core.http.WebSocketClient;
/**
* load and manages a set of HttpClient at boot up.
@@ -84,6 +85,11 @@ public class HttpClients {
}
}
+ public static WebSocketClient createWebSocketClient(HttpClientOptionsSPI
option, boolean sslEnabled) {
+ Vertx vertx = getOrCreateVertx(option);
+ return
vertx.createWebSocketClient(HttpClientOptionsSPI.createWebSocketClientOptions(option,
sslEnabled));
+ }
+
private static Vertx getOrCreateVertx(HttpClientOptionsSPI option) {
if (option.useSharedVertx()) {
return
SharedVertxFactory.getSharedVertx(LegacyPropertyFactory.getEnvironment());
diff --git
a/handlers/handler-loadbalance/src/main/java/org/apache/servicecomb/loadbalance/LoadBalanceFilter.java
b/handlers/handler-loadbalance/src/main/java/org/apache/servicecomb/loadbalance/LoadBalanceFilter.java
index 6f4f1a098..e97263ac4 100644
---
a/handlers/handler-loadbalance/src/main/java/org/apache/servicecomb/loadbalance/LoadBalanceFilter.java
+++
b/handlers/handler-loadbalance/src/main/java/org/apache/servicecomb/loadbalance/LoadBalanceFilter.java
@@ -156,8 +156,11 @@ public class LoadBalanceFilter extends AbstractFilter
implements ConsumerFilter,
}
// user's can invoke a service by supplying target Endpoint.
- // in this case, we do not using load balancer, and no stats of server
calculated, no retrying.
+ // in this case, do not use load balancer, and no stats of server
calculated, no retrying.
private boolean handleSuppliedEndpoint(Invocation invocation) throws
Exception {
+ if (invocation.isEdge()) {
+ return false;
+ }
if (invocation.getEndpoint() != null) {
return true;
}
diff --git
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
b/providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketArgMapperFactory.java
similarity index 54%
copy from
edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
copy to
providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketArgMapperFactory.java
index f7fcfa18c..b0f554a84 100644
---
a/edge/edge-core/src/test/java/org/apache/servicecomb/edge/core/TestUtils.java
+++
b/providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketArgMapperFactory.java
@@ -15,19 +15,22 @@
* limitations under the License.
*/
-package org.apache.servicecomb.edge.core;
+package org.apache.servicecomb.provider.rest.common;
-import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.Test;
+import org.apache.servicecomb.swagger.invocation.arguments.ArgumentMapper;
+import
org.apache.servicecomb.swagger.invocation.arguments.producer.ProducerContextArgumentMapperFactory;
-public class TestUtils {
- @Test
- public void testUtils() {
- Assertions.assertEquals("/a/b/c", Utils.findActualPath("/a/b/c", -1));
- Assertions.assertEquals("/a/b/c", Utils.findActualPath("/a/b/c", 0));
- Assertions.assertEquals("/b/c", Utils.findActualPath("/a/b/c", 1));
- Assertions.assertEquals("/c", Utils.findActualPath("/a/b/c", 2));
- Assertions.assertEquals("", Utils.findActualPath("/a/b/c", 3));
- Assertions.assertEquals("", Utils.findActualPath("/a/b/c", 100));
+import io.vertx.core.http.ServerWebSocket;
+
+public class ProducerServerWebSocketArgMapperFactory implements
ProducerContextArgumentMapperFactory {
+
+ @Override
+ public Class<?> getContextClass() {
+ return ServerWebSocket.class;
+ }
+
+ @Override
+ public ArgumentMapper create(String invocationArgumentName, String
swaggerArgumentName) {
+ return new ProducerServerWebSocketMapper(invocationArgumentName,
swaggerArgumentName);
}
}
diff --git a/core/src/main/java/org/apache/servicecomb/core/Transport.java
b/providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketMapper.java
similarity index 53%
copy from core/src/main/java/org/apache/servicecomb/core/Transport.java
copy to
providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketMapper.java
index 4777a5ca0..21d040b5c 100644
--- a/core/src/main/java/org/apache/servicecomb/core/Transport.java
+++
b/providers/provider-rest-common/src/main/java/org/apache/servicecomb/provider/rest/common/ProducerServerWebSocketMapper.java
@@ -15,38 +15,21 @@
* limitations under the License.
*/
-package org.apache.servicecomb.core;
+package org.apache.servicecomb.provider.rest.common;
-import org.springframework.core.env.Environment;
+import org.apache.servicecomb.common.rest.WebSocketTransportContext;
+import org.apache.servicecomb.swagger.invocation.SwaggerInvocation;
+import
org.apache.servicecomb.swagger.invocation.arguments.producer.AbstractProducerContextArgMapper;
-// TODO:感觉要拆成显式的client、server才好些
-public interface Transport {
- String getName();
- default int getOrder() {
- return 0;
+public class ProducerServerWebSocketMapper extends
AbstractProducerContextArgMapper {
+ public ProducerServerWebSocketMapper(String invocationArgumentName, String
swaggerArgumentName) {
+ super(invocationArgumentName, swaggerArgumentName);
}
- default boolean canInit() {
- return true;
+ @Override
+ public Object createContextArg(SwaggerInvocation invocation) {
+ WebSocketTransportContext context = invocation.getTransportContext();
+ return context.getServerWebSocket();
}
-
- boolean init() throws Exception;
-
- void setEnvironment(Environment environment);
-
- /*
- * endpoint的格式为 URI,比如rest://192.168.1.1:8080
- */
- Object parseAddress(String endpoint);
-
- /*
- * 本transport的监听地址
- */
- Endpoint getEndpoint();
-
- /*
- * 用于上报到服务中心,要求是其他节点可访问的地址
- */
- Endpoint getPublishEndpoint() throws Exception;
}
diff --git
a/providers/provider-rest-common/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.invocation.arguments.producer.ProducerContextArgumentMapperFactory
b/providers/provider-rest-common/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.invocation.arguments.producer.ProducerContextArgumentMapperFactory
index 8898a2e63..1561d349c 100644
---
a/providers/provider-rest-common/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.invocation.arguments.producer.ProducerContextArgumentMapperFactory
+++
b/providers/provider-rest-common/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.invocation.arguments.producer.ProducerContextArgumentMapperFactory
@@ -15,4 +15,5 @@
# limitations under the License.
#
-org.apache.servicecomb.provider.rest.common.ProducerHttpRequestArgMapperFactory
\ No newline at end of file
+org.apache.servicecomb.provider.rest.common.ProducerHttpRequestArgMapperFactory
+org.apache.servicecomb.provider.rest.common.ProducerServerWebSocketArgMapperFactory
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/OperationGenerator.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/OperationGenerator.java
index c7aff4c80..358db8529 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/OperationGenerator.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/OperationGenerator.java
@@ -46,4 +46,10 @@ public interface OperationGenerator {
* Used to check if one of operation form parameter is binary
*/
boolean isBinary();
+
+ /**
+ *
+ * Used to check if this operation is websocket
+ */
+ boolean isWebsocket();
}
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/ParameterGenerator.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/ParameterGenerator.java
index 9b2be38d5..2547b5c7f 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/ParameterGenerator.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/ParameterGenerator.java
@@ -126,7 +126,7 @@ public class ParameterGenerator {
public void generate() {
this.parameterGeneratorContext.updateConsumes(
- this.operationGenerator.isForm(), this.operationGenerator.isBinary());
+ this.operationGenerator.isForm(), this.operationGenerator.isBinary(),
this.operationGenerator.isWebsocket());
if (this.parameterGeneratorContext.getHttpParameterType() ==
HttpParameterType.BODY) {
if (parameterGeneratorContext.getSupportedConsumes().size() == 0) {
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/SwaggerConst.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/SwaggerConst.java
index 5b26bea75..70990922e 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/SwaggerConst.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/SwaggerConst.java
@@ -35,6 +35,10 @@ public final class SwaggerConst {
public static final String PROTOBUF_TYPE = "application/protobuf";
+ public static final String WEBSOCKET_TYPE = "application/websocket";
+
+ public static final String TAG_WEBSOCKET = "websocket";
+
public static final String EXT_JAVA_INTF = "x-java-interface";
public static final String EXT_JAVA_CLASS = "x-java-class";
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/AbstractOperationGenerator.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/AbstractOperationGenerator.java
index d70027593..f3c00478f 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/AbstractOperationGenerator.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/AbstractOperationGenerator.java
@@ -45,6 +45,7 @@ import
org.apache.servicecomb.swagger.generator.ParameterAnnotationProcessor;
import org.apache.servicecomb.swagger.generator.ParameterGenerator;
import org.apache.servicecomb.swagger.generator.ParameterTypeProcessor;
import org.apache.servicecomb.swagger.generator.ResponseTypeProcessor;
+import org.apache.servicecomb.swagger.generator.SwaggerConst;
import org.apache.servicecomb.swagger.generator.SwaggerGeneratorUtils;
import org.apache.servicecomb.swagger.generator.core.model.HttpParameterType;
import org.apache.servicecomb.swagger.generator.core.utils.MethodUtils;
@@ -343,6 +344,12 @@ public abstract class AbstractOperationGenerator
implements OperationGenerator {
return false;
}
+ @Override
+ public boolean isWebsocket() {
+ return this.swaggerOperation.getTags() != null &&
+ this.swaggerOperation.getTags().contains(SwaggerConst.TAG_WEBSOCKET);
+ }
+
@Override
public void addOperationToSwagger() {
if (StringUtils.isEmpty(httpMethod)) {
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/ParameterGeneratorContext.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/ParameterGeneratorContext.java
index 99b1eba22..dc9f6c9c0 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/ParameterGeneratorContext.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/ParameterGeneratorContext.java
@@ -19,6 +19,7 @@ package org.apache.servicecomb.swagger.generator.core;
import java.util.ArrayList;
import java.util.List;
+import org.apache.servicecomb.swagger.generator.SwaggerConst;
import org.apache.servicecomb.swagger.generator.core.model.HttpParameterType;
import com.fasterxml.jackson.databind.JavaType;
@@ -115,8 +116,14 @@ public class ParameterGeneratorContext extends
OperationGeneratorContext {
this.defaultValue = defaultValue;
}
- public void updateConsumes(boolean isForm, boolean isBinary) {
+ public void updateConsumes(boolean isForm, boolean isBinary, boolean
isWebSocket) {
List<String> removed = new ArrayList<>();
+ if (isWebSocket) {
+ supportedConsumes.clear();
+ supportedConsumes.add(SwaggerConst.WEBSOCKET_TYPE);
+ return;
+ }
+
if (httpParameterType == HttpParameterType.BODY) {
if (isForm) {
throw new IllegalArgumentException("Both form and body parameter not
allowed.");
@@ -127,7 +134,11 @@ public class ParameterGeneratorContext extends
OperationGeneratorContext {
}
removed.add(media);
}
- } else if (httpParameterType == HttpParameterType.FORM) {
+ supportedConsumes.removeAll(removed);
+ return;
+ }
+
+ if (httpParameterType == HttpParameterType.FORM) {
for (String media : supportedConsumes) {
if (!SUPPORTED_FORM_CONTENT_TYPE.contains(media)) {
removed.add(media);
@@ -140,10 +151,11 @@ public class ParameterGeneratorContext extends
OperationGeneratorContext {
removed.add(MediaType.MULTIPART_FORM_DATA);
}
}
- } else {
- supportedConsumes.clear();
+ supportedConsumes.removeAll(removed);
+ return;
}
- supportedConsumes.removeAll(removed);
+
+ supportedConsumes.clear();
}
public boolean isForm() {
diff --git
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/SwaggerGeneratorContext.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/SwaggerGeneratorContext.java
index 63d2259e6..7b1eeda13 100644
---
a/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/SwaggerGeneratorContext.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/SwaggerGeneratorContext.java
@@ -30,11 +30,12 @@ import jakarta.ws.rs.core.MediaType;
public class SwaggerGeneratorContext {
protected static final List<String> SUPPORTED_CONTENT_TYPE
= Arrays.asList(MediaType.APPLICATION_JSON, SwaggerConst.PROTOBUF_TYPE,
MediaType.TEXT_PLAIN,
- MediaType.MULTIPART_FORM_DATA, MediaType.APPLICATION_FORM_URLENCODED,
+ MediaType.MULTIPART_FORM_DATA, MediaType.APPLICATION_FORM_URLENCODED,
SwaggerConst.WEBSOCKET_TYPE,
MediaType.SERVER_SENT_EVENTS);
protected static final List<String> SUPPORTED_BODY_CONTENT_TYPE
- = Arrays.asList(MediaType.APPLICATION_JSON, SwaggerConst.PROTOBUF_TYPE,
MediaType.TEXT_PLAIN);
+ = Arrays.asList(MediaType.APPLICATION_JSON, SwaggerConst.PROTOBUF_TYPE,
+ MediaType.TEXT_PLAIN);
protected static final List<String> SUPPORTED_FORM_CONTENT_TYPE
= Arrays.asList(MediaType.MULTIPART_FORM_DATA,
MediaType.APPLICATION_FORM_URLENCODED);
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/processor/parameter/ServerWebSocketContextRegister.java
similarity index 67%
copy from
core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
copy to
swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/processor/parameter/ServerWebSocketContextRegister.java
index a413fce10..70bb42550 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/swagger/swagger-generator/generator-core/src/main/java/org/apache/servicecomb/swagger/generator/core/processor/parameter/ServerWebSocketContextRegister.java
@@ -14,15 +14,18 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.core.invocation;
-import java.util.concurrent.CompletableFuture;
+package org.apache.servicecomb.swagger.generator.core.processor.parameter;
-import org.apache.servicecomb.core.Invocation;
+import java.lang.reflect.Type;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
-public interface InvocationCreator {
- CompletableFuture<Invocation> createAsync();
+import org.apache.servicecomb.swagger.generator.SwaggerContextRegister;
+
+import io.vertx.core.http.ServerWebSocket;
+
+public class ServerWebSocketContextRegister implements SwaggerContextRegister {
+ @Override
+ public Type getContextType() {
+ return ServerWebSocket.class;
+ }
}
diff --git
a/swagger/swagger-generator/generator-core/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.generator.SwaggerContextRegister
b/swagger/swagger-generator/generator-core/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.generator.SwaggerContextRegister
index 55d823147..db52891e0 100644
---
a/swagger/swagger-generator/generator-core/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.generator.SwaggerContextRegister
+++
b/swagger/swagger-generator/generator-core/src/main/resources/META-INF/services/org.apache.servicecomb.swagger.generator.SwaggerContextRegister
@@ -15,4 +15,5 @@
# limitations under the License.
#
-org.apache.servicecomb.swagger.generator.core.processor.parameter.HttpServletRequestContextRegister
\ No newline at end of file
+org.apache.servicecomb.swagger.generator.core.processor.parameter.HttpServletRequestContextRegister
+org.apache.servicecomb.swagger.generator.core.processor.parameter.ServerWebSocketContextRegister
diff --git
a/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/TransportRestClientConfiguration.java
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/TransportRestClientConfiguration.java
index 2b35eb4ca..3775744e1 100644
---
a/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/TransportRestClientConfiguration.java
+++
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/TransportRestClientConfiguration.java
@@ -31,6 +31,11 @@ public class TransportRestClientConfiguration {
return new RestClientCodecFilter();
}
+ @Bean
+ public WebSocketClientCodecFilter webSocketClientCodecFilter() {
+ return new WebSocketClientCodecFilter();
+ }
+
@Bean
public RestClientDecoder restClientDecoder() {
return new RestClientDecoder();
diff --git
a/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientCodecFilter.java
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientCodecFilter.java
new file mode 100644
index 000000000..73a38dfe8
--- /dev/null
+++
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientCodecFilter.java
@@ -0,0 +1,148 @@
+/*
+ * 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.servicecomb.transport.rest.client;
+
+import java.util.concurrent.CompletableFuture;
+
+import org.apache.commons.lang3.StringUtils;
+import org.apache.servicecomb.common.rest.RestConst;
+import org.apache.servicecomb.common.rest.WebSocketTransportContext;
+import org.apache.servicecomb.common.rest.definition.RestMetaUtils;
+import org.apache.servicecomb.common.rest.definition.RestOperationMeta;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.filter.AbstractFilter;
+import org.apache.servicecomb.core.filter.ConsumerFilter;
+import org.apache.servicecomb.core.filter.EdgeFilter;
+import org.apache.servicecomb.core.filter.Filter;
+import org.apache.servicecomb.core.filter.FilterNode;
+import org.apache.servicecomb.foundation.common.net.URIEndpointObject;
+import org.apache.servicecomb.foundation.common.utils.SPIServiceUtils;
+import
org.apache.servicecomb.foundation.vertx.client.http.HttpClientOptionsSPI;
+import org.apache.servicecomb.foundation.vertx.client.http.HttpClients;
+import org.apache.servicecomb.registry.definition.DefinitionConst;
+import org.apache.servicecomb.swagger.invocation.Response;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.vertx.core.http.ServerWebSocket;
+import io.vertx.core.http.WebSocket;
+import io.vertx.core.http.WebSocketClient;
+
+public class WebSocketClientCodecFilter extends AbstractFilter implements
ConsumerFilter, EdgeFilter {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(WebSocketClientCodecFilter.class);
+
+ public static final String NAME = "websocket-client";
+
+ @Override
+ public String getName() {
+ return NAME;
+ }
+
+ @Override
+ public boolean enabledForTransport(String transport) {
+ return CoreConst.WEBSOCKET.equals(transport);
+ }
+
+ @Override
+ public int getOrder() {
+ return Filter.CONSUMER_LOAD_BALANCE_ORDER + 2000;
+ }
+
+ @Override
+ public CompletableFuture<Response> onFilter(Invocation invocation,
FilterNode nextNode) {
+ invocation.getInvocationStageTrace().startConsumerConnection();
+
+ CompletableFuture<Response> createWebSocket = new CompletableFuture<>();
+ URIEndpointObject endpoint = (URIEndpointObject)
invocation.getEndpoint().getAddress();
+ HttpClientOptionsSPI optionsSPI;
+ if (endpoint.isHttp2Enabled()) {
+ optionsSPI = SPIServiceUtils.getTargetService(HttpClientOptionsSPI.class,
+ HttpTransportHttpClientOptionsSPI.class);
+ } else {
+ optionsSPI = SPIServiceUtils.getTargetService(HttpClientOptionsSPI.class,
+ Http2TransportHttpClientOptionsSPI.class);
+ }
+ WebSocketClient webSocketClient =
HttpClients.createWebSocketClient(optionsSPI, endpoint.isSslEnabled());
+
+ try {
+ webSocketClient.connect(endpoint.getPort(), endpoint.getHostOrIp(),
createRequestPath(invocation,
+
RestMetaUtils.getRestOperationMeta(invocation.getOperationMeta())))
+ .onComplete(asyncResult -> {
+ invocation.getInvocationStageTrace().finishConsumerConnection();
+ if (asyncResult.failed()) {
+ createWebSocket.completeExceptionally(asyncResult.cause());
+ return;
+ }
+ if (invocation.isEdge()) {
+ WebSocketTransportContext parentContext =
invocation.getTransportContext();
+ ServerWebSocket serverWebSocket =
parentContext.getServerWebSocket();
+ WebSocket clientWebSocket = asyncResult.result();
+ serverWebSocket.closeHandler(v -> {
+ if (!clientWebSocket.isClosed()) {
+ clientWebSocket.close();
+ }
+ });
+
serverWebSocket.textMessageHandler(clientWebSocket::writeTextMessage);
+
serverWebSocket.binaryMessageHandler(clientWebSocket::writeBinaryMessage);
+ serverWebSocket.exceptionHandler(e -> {
+ LOGGER.warn("consumer exception.", e);
+ if (!serverWebSocket.isClosed()) {
+ serverWebSocket.close();
+ }
+ });
+ clientWebSocket.closeHandler(v -> {
+ if (!serverWebSocket.isClosed()) {
+ serverWebSocket.close();
+ }
+ });
+
clientWebSocket.textMessageHandler(serverWebSocket::writeTextMessage);
+
clientWebSocket.binaryMessageHandler(serverWebSocket::writeBinaryMessage);
+ clientWebSocket.exceptionHandler(e -> {
+ LOGGER.warn("producer exception.", e);
+ if (!clientWebSocket.isClosed()) {
+ clientWebSocket.close();
+ }
+ });
+ }
+ invocation.setTransportContext(new WebSocketClientTransportContext(
+ asyncResult.result()));
+
createWebSocket.complete(Response.createSuccess(asyncResult.result()));
+ });
+ } catch (Exception e) {
+ createWebSocket.completeExceptionally(e);
+ }
+
+ return createWebSocket;
+ }
+
+ protected String createRequestPath(Invocation invocation, RestOperationMeta
restOperationMeta) throws Exception {
+ String path =
invocation.getLocalContext(RestConst.REST_CLIENT_REQUEST_PATH);
+ if (path == null) {
+ path =
restOperationMeta.getPathBuilder().createRequestPath(invocation.getSwaggerArguments());
+ }
+
+ URIEndpointObject endpoint = (URIEndpointObject)
invocation.getEndpoint().getAddress();
+ String urlPrefix = endpoint.getFirst(DefinitionConst.URL_PREFIX);
+ if (StringUtils.isEmpty(urlPrefix) || path.startsWith(urlPrefix)) {
+ return path;
+ }
+
+ return urlPrefix + path;
+ }
+}
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientTransportContext.java
similarity index 62%
copy from
core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
copy to
transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientTransportContext.java
index a413fce10..89d43bcde 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/transports/transport-rest/transport-rest-client/src/main/java/org/apache/servicecomb/transport/rest/client/WebSocketClientTransportContext.java
@@ -14,15 +14,22 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.core.invocation;
-import java.util.concurrent.CompletableFuture;
+package org.apache.servicecomb.transport.rest.client;
-import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.swagger.invocation.context.TransportContext;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
-public interface InvocationCreator {
- CompletableFuture<Invocation> createAsync();
+import io.vertx.core.http.WebSocket;
+
+public class WebSocketClientTransportContext implements TransportContext {
+
+ protected final WebSocket webSocketClient;
+
+ public WebSocketClientTransportContext(WebSocket webSocketClient) {
+ this.webSocketClient = webSocketClient;
+ }
+
+ public WebSocket getWebSocketClient() {
+ return webSocketClient;
+ }
}
diff --git
a/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/RestServerVerticle.java
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/RestServerVerticle.java
index c72e50435..e6ecb2804 100644
---
a/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/RestServerVerticle.java
+++
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/RestServerVerticle.java
@@ -69,11 +69,14 @@ public class RestServerVerticle extends AbstractVerticle {
private URIEndpointObject endpointObject;
+ private WebSocketDispatcher webSocketDispatcher;
+
@Override
public void init(Vertx vertx, Context context) {
super.init(vertx, context);
Endpoint endpoint = (Endpoint)
context.config().getValue(AbstractTransport.ENDPOINT_KEY);
this.endpointObject = (URIEndpointObject) endpoint.getAddress();
+ this.webSocketDispatcher = new WebSocketDispatcher(endpoint);
}
@Override
@@ -92,7 +95,17 @@ public class RestServerVerticle extends AbstractVerticle {
initDispatcher(mainRouter);
mountGlobalRestFailureHandler(mainRouter);
HttpServer httpServer = createHttpServer();
- httpServer.requestHandler(mainRouter);
+ httpServer.requestHandler(httpServerRequest -> {
+ if (this.endpointObject.isWebsocketEnabled()) {
+ String path =
LegacyPropertyFactory.getStringProperty("servicecomb.rest.server.websocket-prefix");
+ if (httpServerRequest.path().startsWith(path)) {
+ httpServerRequest.toWebSocket().onComplete(w ->
webSocketDispatcher.onRequest(w),
+ e -> LOGGER.error("WebSocket error.", e));
+ return;
+ }
+ }
+ mainRouter.handle(httpServerRequest);
+ });
httpServer.connectionHandler(connection -> {
DefaultHttpServerMetrics serverMetrics = (DefaultHttpServerMetrics)
((ConnectionBase) connection).metrics();
DefaultServerEndpointMetric endpointMetric =
serverMetrics.getEndpointMetric();
diff --git
a/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketDispatcher.java
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketDispatcher.java
new file mode 100644
index 000000000..5f8e2cd32
--- /dev/null
+++
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketDispatcher.java
@@ -0,0 +1,122 @@
+/*
+ * 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.servicecomb.transport.rest.vertx;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import org.apache.servicecomb.common.rest.EdgeServerWebSocketInvocationCreator;
+import
org.apache.servicecomb.common.rest.ProviderServerWebSocketInvocationCreator;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationItem;
+import org.apache.servicecomb.common.rest.route.URLMappedConfigurationLoader;
+import org.apache.servicecomb.common.rest.route.Utils;
+import org.apache.servicecomb.config.ConfigurationChangedEvent;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.Endpoint;
+import org.apache.servicecomb.core.SCBEngine;
+import org.apache.servicecomb.core.Transport;
+import org.apache.servicecomb.core.definition.MicroserviceMeta;
+import org.apache.servicecomb.core.invocation.InvocationCreator;
+import org.apache.servicecomb.foundation.common.LegacyPropertyFactory;
+import org.apache.servicecomb.foundation.common.event.EventManager;
+import org.apache.servicecomb.swagger.invocation.exception.CommonExceptionData;
+import org.apache.servicecomb.swagger.invocation.exception.InvocationException;
+
+import com.google.common.eventbus.Subscribe;
+
+import io.vertx.core.http.ServerWebSocket;
+import jakarta.ws.rs.core.Response.Status;
+
+public class WebSocketDispatcher {
+ private static final String KEY_MAPPING_PREFIX =
"servicecomb.http.dispatcher.edge.websocket.mappings";
+
+ private final Object LOCK = new Object();
+
+ private volatile boolean initialized = false;
+
+ private Endpoint endpoint;
+
+ private MicroserviceMeta microserviceMeta;
+
+ private boolean isEdge;
+
+ private Map<String, URLMappedConfigurationItem> configurations = new
HashMap<>();
+
+ public WebSocketDispatcher(Endpoint endpoint) {
+ this.endpoint = endpoint;
+ EventManager.register(this);
+ }
+
+ private void loadConfigurations() {
+ configurations = URLMappedConfigurationLoader.loadConfigurations(
+ LegacyPropertyFactory.getEnvironment(), KEY_MAPPING_PREFIX);
+ }
+
+ @Subscribe
+ public void onConfigurationChangedEvent(ConfigurationChangedEvent event) {
+ for (String changed : event.getChanged()) {
+ if (changed.startsWith(KEY_MAPPING_PREFIX)) {
+ loadConfigurations();
+ break;
+ }
+ }
+ }
+
+ protected void onRequest(ServerWebSocket webSocket) {
+ if (!initialized) {
+ synchronized (LOCK) {
+ if (!initialized) {
+ Transport transport =
SCBEngine.getInstance().getTransportManager().findTransport(CoreConst.WEBSOCKET);
+ this.microserviceMeta =
SCBEngine.getInstance().getProducerMicroserviceMeta();
+ this.endpoint = new Endpoint(transport, this.endpoint.getEndpoint());
+ this.isEdge = TransportConfig.getRestServerVerticle()
+
.getName().equals("org.apache.servicecomb.edge.core.EdgeRestServerVerticle");
+ if (this.isEdge) {
+ loadConfigurations();
+ }
+ }
+ initialized = true;
+ }
+ }
+
+ InvocationCreator creator;
+ if (isEdge) {
+ URLMappedConfigurationItem configurationItem =
findConfigurationItem(webSocket.path());
+ if (configurationItem == null) {
+ throw new InvocationException(Status.NOT_FOUND, new
CommonExceptionData(
+ String.format("path %s not found", webSocket.path())));
+ }
+ String path = Utils.findActualPath(webSocket.path(),
configurationItem.getPrefixSegmentCount());
+ creator = new EdgeServerWebSocketInvocationCreator(
+ configurationItem.getMicroserviceName(), path, endpoint, webSocket);
+ } else {
+ creator = new ProviderServerWebSocketInvocationCreator(microserviceMeta,
+ endpoint, webSocket);
+ }
+ new WebSocketProducerInvocationFlow(creator, webSocket).run();
+ }
+
+ private URLMappedConfigurationItem findConfigurationItem(String path) {
+ for (URLMappedConfigurationItem item : configurations.values()) {
+ if (item.getPattern().matcher(path).matches()) {
+ return item;
+ }
+ }
+ return null;
+ }
+}
diff --git
a/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketProducerInvocationFlow.java
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketProducerInvocationFlow.java
new file mode 100644
index 000000000..1297da53a
--- /dev/null
+++
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketProducerInvocationFlow.java
@@ -0,0 +1,51 @@
+/*
+ * 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.servicecomb.transport.rest.vertx;
+
+import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.invocation.InvocationCreator;
+import org.apache.servicecomb.core.invocation.ProducerInvocationFlow;
+import org.apache.servicecomb.swagger.invocation.Response;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.vertx.core.http.ServerWebSocket;
+
+public class WebSocketProducerInvocationFlow extends ProducerInvocationFlow {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(WebSocketProducerInvocationFlow.class);
+
+ private final ServerWebSocket websocket;
+
+ public WebSocketProducerInvocationFlow(InvocationCreator invocationCreator,
ServerWebSocket webSocket) {
+ super(invocationCreator);
+ this.websocket = webSocket;
+ }
+
+ @Override
+ protected Invocation sendCreateInvocationException(Throwable throwable) {
+ LOGGER.error("Web socket create invocation error.", throwable);
+ websocket.writeTextMessage("Web socket create invocation error " +
throwable.getMessage());
+ websocket.close();
+ return null;
+ }
+
+ @Override
+ protected void endResponse(Invocation invocation, Response response) {
+
+ }
+}
diff --git
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketTransport.java
similarity index 64%
copy from
core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
copy to
transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketTransport.java
index a413fce10..9bcd614a3 100644
---
a/core/src/main/java/org/apache/servicecomb/core/invocation/InvocationCreator.java
+++
b/transports/transport-rest/transport-rest-vertx/src/main/java/org/apache/servicecomb/transport/rest/vertx/WebSocketTransport.java
@@ -14,15 +14,25 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.servicecomb.core.invocation;
-import java.util.concurrent.CompletableFuture;
+package org.apache.servicecomb.transport.rest.vertx;
-import org.apache.servicecomb.core.Invocation;
+import org.apache.servicecomb.core.CoreConst;
+import org.apache.servicecomb.core.transport.AbstractTransport;
-/**
- * better to named InvocationFactory, but already be used by old version
- */
-public interface InvocationCreator {
- CompletableFuture<Invocation> createAsync();
+public class WebSocketTransport extends AbstractTransport {
+ @Override
+ public String getName() {
+ return CoreConst.WEBSOCKET;
+ }
+
+ @Override
+ public int getOrder() {
+ return -500;
+ }
+
+ @Override
+ public boolean init() throws Exception {
+ return true;
+ }
}
diff --git
a/transports/transport-rest/transport-rest-vertx/src/main/resources/META-INF/services/org.apache.servicecomb.core.Transport
b/transports/transport-rest/transport-rest-vertx/src/main/resources/META-INF/services/org.apache.servicecomb.core.Transport
index be04a13de..74ae694b9 100644
---
a/transports/transport-rest/transport-rest-vertx/src/main/resources/META-INF/services/org.apache.servicecomb.core.Transport
+++
b/transports/transport-rest/transport-rest-vertx/src/main/resources/META-INF/services/org.apache.servicecomb.core.Transport
@@ -16,3 +16,4 @@
#
org.apache.servicecomb.transport.rest.vertx.VertxRestTransport
+org.apache.servicecomb.transport.rest.vertx.WebSocketTransport