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 71279a27e [SCB-2885]support write SSE response (#4383)
71279a27e is described below
commit 71279a27e0340c86f741eb341ba3ec0612b29bf3
Author: liubao68 <[email protected]>
AuthorDate: Wed Jun 26 11:37:44 2024 +0800
[SCB-2885]support write SSE response (#4383)
---
.../codec/produce/ProduceEventStreamProcessor.java | 56 ++++++++++++++++++++++
.../codec/produce/ProduceProcessorManager.java | 42 ++++++++++------
.../rest/filter/inner/RestServerCodecFilter.java | 5 +-
.../samples/ReactiveStreamController.java | 6 ++-
.../provider/src/main/resources/application.yml | 6 +++
5 files changed, 97 insertions(+), 18 deletions(-)
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceEventStreamProcessor.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceEventStreamProcessor.java
new file mode 100644
index 000000000..3491539d1
--- /dev/null
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceEventStreamProcessor.java
@@ -0,0 +1,56 @@
+/*
+ * 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.codec.produce;
+
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.nio.charset.StandardCharsets;
+
+import org.apache.servicecomb.common.rest.codec.RestObjectMapperFactory;
+
+import com.fasterxml.jackson.databind.JavaType;
+
+import jakarta.ws.rs.core.MediaType;
+
+public class ProduceEventStreamProcessor implements ProduceProcessor {
+ private int writeIndex = 0;
+
+ @Override
+ public String getName() {
+ return MediaType.SERVER_SENT_EVENTS;
+ }
+
+ @Override
+ public int getOrder() {
+ return 0;
+ }
+
+ @Override
+ public void doEncodeResponse(OutputStream output, Object result) throws
Exception {
+ String buffer = "id: " + (writeIndex++) + "\n"
+ + "data: "
+ +
RestObjectMapperFactory.getRestObjectMapper().writeValueAsString(result)
+ + "\n\n";
+ output.write(buffer.getBytes(StandardCharsets.UTF_8));
+ }
+
+ @Override
+ public Object doDecodeResponse(InputStream input, JavaType type) throws
Exception {
+ return null;
+ }
+}
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceProcessorManager.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceProcessorManager.java
index 7182aff20..3fdcfec13 100644
---
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceProcessorManager.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/codec/produce/ProduceProcessorManager.java
@@ -120,32 +120,43 @@ public final class ProduceProcessorManager extends
RegisterManager<String, Map<S
public ProduceProcessor createProduceProcessor(OperationMeta operationMeta,
int statusCode, String accept, Class<?> serialViewClass) {
+ // If no produces defined, using default processor
ApiResponses responses =
operationMeta.getSwaggerOperation().getResponses();
ApiResponse response = responses.get(String.valueOf(statusCode));
if (response == null || response.getContent() == null ||
response.getContent().size() == 0) {
return findDefaultProcessor();
}
- String actualAccept = accept;
- if (actualAccept == null) {
+
+ // check intersection of `Accept` and `Produces`
+ if (accept == null) {
if (response.getContent().get(defaultResponseEncoding()) != null) {
- actualAccept = defaultResponseEncoding();
+ accept = defaultResponseEncoding();
} else {
- actualAccept = response.getContent().keySet().iterator().next();
+ accept = response.getContent().keySet().iterator().next();
}
}
- ContentType contentType = ContentType.parse(actualAccept);
- actualAccept = contentType.getMimeType();
- if (MediaType.WILDCARD.equals(contentType.getMimeType()) ||
- MediaType.MEDIA_TYPE_WILDCARD.equals(contentType.getMimeType())) {
- if (response.getContent().get(defaultResponseEncoding()) != null) {
- actualAccept = defaultResponseEncoding();
- } else {
- actualAccept = response.getContent().keySet().iterator().next();
+
+ String actualAccept = null;
+ for (String item : accept.split(",")) {
+ ContentType contentType = ContentType.parse(item);
+ if (MediaType.WILDCARD.equals(contentType.getMimeType()) ||
+ MediaType.MEDIA_TYPE_WILDCARD.equals(contentType.getMimeType())) {
+ if (response.getContent().get(defaultResponseEncoding()) != null) {
+ actualAccept = defaultResponseEncoding();
+ } else {
+ actualAccept = response.getContent().keySet().iterator().next();
+ }
+ break;
+ }
+ if (response.getContent().get(contentType.getMimeType()) != null) {
+ actualAccept = contentType.getMimeType();
+ break;
}
}
- if (response.getContent().get(actualAccept) == null) {
- LOGGER.warn("Operation do not support accept type {}/{}", accept,
actualAccept);
+
+ if (actualAccept == null) {
+ LOGGER.warn("Operation {} do not support accept type {}",
operationMeta.getSchemaQualifiedName(), accept);
return findDefaultProcessor();
}
if (MediaType.APPLICATION_JSON.equals(actualAccept)) {
@@ -155,6 +166,9 @@ public final class ProduceProcessorManager extends
RegisterManager<String, Map<S
return new ProduceProtoBufferProcessor(operationMeta,
operationMeta.getSchemaMeta().getSwagger(),
response.getContent().get(actualAccept).getSchema());
}
+ if (MediaType.SERVER_SENT_EVENTS.equals(actualAccept)) {
+ return new ProduceEventStreamProcessor();
+ }
// text plain
return findPlainProcessorByViewClass(serialViewClass);
}
diff --git
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/RestServerCodecFilter.java
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/RestServerCodecFilter.java
index d47dd4bce..c46a47c1f 100644
---
a/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/RestServerCodecFilter.java
+++
b/common/common-rest/src/main/java/org/apache/servicecomb/common/rest/filter/inner/RestServerCodecFilter.java
@@ -180,14 +180,14 @@ public class RestServerCodecFilter extends AbstractFilter
implements ProviderFil
if (isServerSendEvent(response)) {
responseEx.setContentType(produceProcessor.getName());
- writeServerSendEvent(response, produceProcessor, responseEx);
+ return writeServerSendEvent(response, produceProcessor, responseEx);
}
responseEx.setContentType(produceProcessor.getName());
return writeResponse(responseEx, produceProcessor, response.getResult(),
response, true);
}
- private static void writeServerSendEvent(Response response, ProduceProcessor
produceProcessor,
+ private static CompletableFuture<Response> writeServerSendEvent(Response
response, ProduceProcessor produceProcessor,
HttpServletResponseEx responseEx) {
responseEx.setChunked(true);
CompletableFuture<Response> result = new CompletableFuture<>();
@@ -223,6 +223,7 @@ public class RestServerCodecFilter extends AbstractFilter
implements ProviderFil
result.complete(response);
}
});
+ return result;
}
/**
diff --git
a/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/ReactiveStreamController.java
b/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/ReactiveStreamController.java
index b72187b09..33458737f 100644
---
a/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/ReactiveStreamController.java
+++
b/demo/demo-zookeeper/provider/src/main/java/org/apache/servicecomb/samples/ReactiveStreamController.java
@@ -17,6 +17,8 @@
package org.apache.servicecomb.samples;
+import java.util.concurrent.TimeUnit;
+
import org.apache.servicecomb.provider.rest.common.RestSchema;
import org.reactivestreams.Publisher;
import org.springframework.web.bind.annotation.GetMapping;
@@ -67,7 +69,7 @@ public class ReactiveStreamController {
@GetMapping("/sseModel")
public Publisher<Model> sseModel() {
- return Flowable.fromArray(new Model("a", 1), new Model("b", 2),
- new Model("c", 3));
+ return Flowable.intervalRange(0, 5, 0, 3, TimeUnit.SECONDS)
+ .map(item -> new Model("jack", item.intValue()));
}
}
diff --git a/demo/demo-zookeeper/provider/src/main/resources/application.yml
b/demo/demo-zookeeper/provider/src/main/resources/application.yml
index 5707f0c12..c8a158d31 100644
--- a/demo/demo-zookeeper/provider/src/main/resources/application.yml
+++ b/demo/demo-zookeeper/provider/src/main/resources/application.yml
@@ -34,3 +34,9 @@ servicecomb:
rest:
address: 0.0.0.0:9094
+ cors:
+ enabled: true
+ origin: "*"
+ allowCredentials: false
+ allowedMethod: "*"
+ maxAge: 3600