This is an automated email from the ASF dual-hosted git repository.

szetszwo pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ratis.git


The following commit(s) were added to refs/heads/master by this push:
     new 44152380d RATIS-2604. Add a new server side DataStreamApi (#1518)
44152380d is described below

commit 44152380d5ec1c904824c0cdf0ec853c2a3486c8
Author: Peter Lee <[email protected]>
AuthorDate: Tue Jul 14 18:59:04 2026 +0800

    RATIS-2604. Add a new server side DataStreamApi (#1518)
---
 .../examples/filestore/FileStoreStateMachine.java  |  9 ++++--
 .../ratis/netty/server/ReadStreamManagement.java   |  2 +-
 .../org/apache/ratis/server/api/DataStreamApi.java | 37 ++++++++++++++++++++++
 .../apache/ratis/statemachine/StateMachine.java    |  8 +++--
 .../ratis/datastream/DataStreamTestUtils.java      | 17 ++++------
 .../netty/server/TestDataStreamManagement.java     | 12 ++++---
 6 files changed, 64 insertions(+), 21 deletions(-)

diff --git 
a/ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileStoreStateMachine.java
 
b/ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileStoreStateMachine.java
index 4ff83369a..5e9bc07e0 100644
--- 
a/ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileStoreStateMachine.java
+++ 
b/ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileStoreStateMachine.java
@@ -108,20 +108,23 @@ public class FileStoreStateMachine extends 
BaseStateMachine {
   }
 
   @Override
-  public void query(Message request, WritableByteChannel stream) {
+  public long transferTo(Message request, WritableByteChannel stream) throws 
IOException {
     try {
       final ReadRequestProto proto = 
ReadRequestProto.parseFrom(request.getContent());
       if (proto.getIsWatch()) {
         throw new IOException("Watch is not supported for streaming read: " + 
proto);
       }
-      files.streamRead(proto.getPath().toStringUtf8(), proto.getOffset(), 
proto.getLength(), stream);
-    } catch (Exception e) {
+      final long length = proto.getLength();
+      files.streamRead(proto.getPath().toStringUtf8(), proto.getOffset(), 
length, stream);
+      return length;
+    } catch (IOException e) {
       LOG.error(getId() + ": Failed streaming read for " + request, e);
       try {
         stream.close();
       } catch (IOException ignored) {
         // ignore
       }
+      throw e;
     }
   }
 
diff --git 
a/ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java
 
b/ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java
index 46d169bbf..7f3796790 100644
--- 
a/ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java
+++ 
b/ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java
@@ -211,7 +211,7 @@ public class ReadStreamManagement {
 
       final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx, readCheckReply);
       try {
-        division.getStateMachine().data().query(request.getMessage(), stream);
+        division.getStateMachine().data().transferTo(request.getMessage(), 
stream);
       } catch (Throwable t) {
         LOG.error("{}: Failed read-only data stream query for {}", this, 
request, t);
       }
diff --git 
a/ratis-server-api/src/main/java/org/apache/ratis/server/api/DataStreamApi.java 
b/ratis-server-api/src/main/java/org/apache/ratis/server/api/DataStreamApi.java
new file mode 100644
index 000000000..bf5335063
--- /dev/null
+++ 
b/ratis-server-api/src/main/java/org/apache/ratis/server/api/DataStreamApi.java
@@ -0,0 +1,37 @@
+/*
+ * 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.ratis.server.api;
+
+import org.apache.ratis.protocol.Message;
+
+import java.io.IOException;
+import java.nio.channels.WritableByteChannel;
+
+/** Server side data stream API. */
+public interface DataStreamApi {
+  /**
+   * Similar to {@link java.nio.channels.FileChannel#transferTo(long, long, 
WritableByteChannel)}
+   * except that the parameter of this method is request message
+   * but not position, count in FileChannel.
+   *
+   * @param request the request message
+   * @param stream the output stream to send the results
+   * @return the number of bytes transferred
+   */
+  long transferTo(Message request, WritableByteChannel stream) throws 
IOException;
+}
diff --git 
a/ratis-server-api/src/main/java/org/apache/ratis/statemachine/StateMachine.java
 
b/ratis-server-api/src/main/java/org/apache/ratis/statemachine/StateMachine.java
index d2443d9b5..50b7d7f6d 100644
--- 
a/ratis-server-api/src/main/java/org/apache/ratis/statemachine/StateMachine.java
+++ 
b/ratis-server-api/src/main/java/org/apache/ratis/statemachine/StateMachine.java
@@ -26,6 +26,7 @@ import org.apache.ratis.protocol.RaftGroupMemberId;
 import org.apache.ratis.protocol.RaftPeer;
 import org.apache.ratis.protocol.RaftPeerId;
 import org.apache.ratis.server.RaftServer;
+import org.apache.ratis.server.api.DataStreamApi;
 import org.apache.ratis.server.protocol.TermIndex;
 import org.apache.ratis.server.storage.RaftStorage;
 import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
@@ -66,7 +67,7 @@ public interface StateMachine extends Closeable {
    * For data intensive applications, it can be more efficient to implement 
this API
    * in order to support zero buffer coping and a light-weighted raft log.
    */
-  interface DataApi {
+  interface DataApi extends DataStreamApi {
     /** A noop implementation of {@link DataApi}. */
     DataApi DEFAULT = new DataApi() {};
 
@@ -123,8 +124,11 @@ public interface StateMachine extends Closeable {
      *
      * @param request the client request
      * @param stream the output stream to send the results
+     * @return the number of bytes transferred
      */
-    default void query(Message request, WritableByteChannel stream) {
+    @Override
+    default long transferTo(Message request, WritableByteChannel stream) 
throws IOException {
+      return 0;
     }
 
     /**
diff --git 
a/ratis-test/src/test/java/org/apache/ratis/datastream/DataStreamTestUtils.java 
b/ratis-test/src/test/java/org/apache/ratis/datastream/DataStreamTestUtils.java
index fe9c3f9ea..dcf2c60ac 100644
--- 
a/ratis-test/src/test/java/org/apache/ratis/datastream/DataStreamTestUtils.java
+++ 
b/ratis-test/src/test/java/org/apache/ratis/datastream/DataStreamTestUtils.java
@@ -190,26 +190,21 @@ public interface DataStreamTestUtils {
     }
 
     @Override
-    public void query(Message request, WritableByteChannel stream) {
-      CompletableFuture.supplyAsync(() -> {
-        try {
-          streamReadOnlyImpl(request, stream);
-        } catch (IOException e) {
-          throw new CompletionException("Failed to streamReadOnly for " + 
request, e);
-        }
-        return null;
-      });
+    public long transferTo(Message request, WritableByteChannel stream) throws 
IOException {
+      return streamReadOnlyImpl(request, stream);
     }
 
-    private void streamReadOnlyImpl(Message request, WritableByteChannel 
stream) throws IOException {
+    private long streamReadOnlyImpl(Message request, WritableByteChannel 
stream) throws IOException {
+      long transferred = 0;
       try {
         for (int i = 0; i < READ_ONLY_STREAM_CHUNKS; i++) {
           final ByteString chunk = 
getReadOnlyStreamChunk(request.getContent(), i);
-          stream.write(chunk.asReadOnlyByteBuffer());
+          transferred += stream.write(chunk.asReadOnlyByteBuffer());
         }
       } finally {
         stream.close();
       }
+      return transferred;
     }
 
     SingleDataStream getSingleDataStream(RaftClientRequest request) {
diff --git 
a/ratis-test/src/test/java/org/apache/ratis/netty/server/TestDataStreamManagement.java
 
b/ratis-test/src/test/java/org/apache/ratis/netty/server/TestDataStreamManagement.java
index 56ecd166b..4d311706a 100644
--- 
a/ratis-test/src/test/java/org/apache/ratis/netty/server/TestDataStreamManagement.java
+++ 
b/ratis-test/src/test/java/org/apache/ratis/netty/server/TestDataStreamManagement.java
@@ -88,9 +88,10 @@ class TestDataStreamManagement {
     final AtomicReference<WritableByteChannel> streamRef = new 
AtomicReference<>();
     final DataApi dataApi = new DataApi() {
       @Override
-      public void query(Message request, WritableByteChannel stream) {
+      public long transferTo(Message request, WritableByteChannel stream) {
         messageRef.set(request);
         streamRef.set(stream);
+        return 0;
       }
     };
     final ReadStreamManagement management = newReadStreamManagement(serverId, 
groupId, dataApi);
@@ -140,9 +141,10 @@ class TestDataStreamManagement {
 
     final DataApi dataApi = new DataApi() {
       @Override
-      public void query(Message request, WritableByteChannel stream) {
+      public long transferTo(Message request, WritableByteChannel stream) {
         messageRef.set(request);
         streamRef.set(stream);
+        return 0;
       }
     };
     final ReadStreamManagement management = newReadStreamManagement(serverId, 
groupId, dataApi, request -> {
@@ -181,8 +183,9 @@ class TestDataStreamManagement {
 
     final DataApi dataApi = new DataApi() {
       @Override
-      public void query(Message request, WritableByteChannel stream) {
+      public long transferTo(Message request, WritableByteChannel stream) {
         queryCalled.set(true);
+        return 0;
       }
     };
     final ReadStreamManagement management = newReadStreamManagement(serverId, 
groupId, dataApi, request ->
@@ -234,9 +237,10 @@ class TestDataStreamManagement {
 
     final DataApi dataApi = new DataApi() {
       @Override
-      public void query(Message request, WritableByteChannel stream) {
+      public long transferTo(Message request, WritableByteChannel stream) {
         queryInEventLoop.set(eventLoop.inEventLoop());
         queryDone.countDown();
+        return 0;
       }
     };
     final StateMachine stateMachine = new BaseStateMachine() {

Reply via email to