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

Reply via email to