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 27cdb5d38 RATIS-2603. Add DataStreamApi.Resolver for group-independent 
read-only requests (#1516)
27cdb5d38 is described below

commit 27cdb5d3816ac1d8bc69e734e0fec8037ca1acda
Author: Peter Lee <[email protected]>
AuthorDate: Wed Jul 15 01:32:00 2026 +0800

    RATIS-2603. Add DataStreamApi.Resolver for group-independent read-only 
requests (#1516)
---
 .../ratis/netty/server/NettyServerStreamRpc.java   |   2 +-
 .../ratis/netty/server/ReadStreamManagement.java   |  62 +++++++-
 .../apache/ratis/server/RaftServerConfigKeys.java  |  13 ++
 .../org/apache/ratis/server/api/DataStreamApi.java |   8 +
 .../netty/server/TestDataStreamManagement.java     | 171 +++++++++++++++++++++
 5 files changed, 248 insertions(+), 8 deletions(-)

diff --git 
a/ratis-netty/src/main/java/org/apache/ratis/netty/server/NettyServerStreamRpc.java
 
b/ratis-netty/src/main/java/org/apache/ratis/netty/server/NettyServerStreamRpc.java
index d7cdeb056..b9bbb23bc 100644
--- 
a/ratis-netty/src/main/java/org/apache/ratis/netty/server/NettyServerStreamRpc.java
+++ 
b/ratis-netty/src/main/java/org/apache/ratis/netty/server/NettyServerStreamRpc.java
@@ -165,7 +165,7 @@ public class NettyServerStreamRpc implements 
DataStreamServerRpc {
     this.name = server.getId() + "-" + 
JavaUtils.getClassSimpleName(getClass());
     this.metrics = new NettyServerStreamRpcMetrics(this.name);
     this.requests = new DataStreamManagement(server, metrics);
-    this.reads = new ReadStreamManagement(server);
+    this.reads = new ReadStreamManagement(server, 
RaftServerConfigKeys.DataStream.serverApiResolver(parameters));
 
     final RaftProperties properties = server.getProperties();
 
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 7f3796790..9a3b51caf 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
@@ -29,6 +29,8 @@ import org.apache.ratis.protocol.ClientId;
 import org.apache.ratis.protocol.RaftClientReply;
 import org.apache.ratis.protocol.RaftClientRequest;
 import org.apache.ratis.protocol.exceptions.AlreadyClosedException;
+import org.apache.ratis.protocol.exceptions.DataStreamException;
+import org.apache.ratis.server.api.DataStreamApi;
 import org.apache.ratis.server.RaftServer;
 import org.apache.ratis.server.RaftServerConfigKeys;
 import 
org.apache.ratis.thirdparty.com.google.protobuf.InvalidProtocolBufferException;
@@ -91,8 +93,16 @@ public class ReadStreamManagement {
 
     @Override
     public synchronized void close() throws IOException {
+      close(terminalReply);
+    }
+
+    private synchronized void fail(RaftClientReply reply) throws IOException {
+      close(newReadStreamTerminalReply(clientId, streamId, reply));
+    }
+
+    private void close(DataStreamReplyByteBuffer reply) throws IOException {
       if (closed.complete(null)) {
-        writeAndFlush(terminalReply, 0);
+        writeAndFlush(reply, 0);
       }
     }
 
@@ -137,11 +147,17 @@ public class ReadStreamManagement {
   }
 
   private final RaftServer server;
+  private final DataStreamApi.Resolver readResolver;
   private final String name;
   private final ExecutorService requestExecutor;
 
   ReadStreamManagement(RaftServer server) {
+    this(server, null);
+  }
+
+  ReadStreamManagement(RaftServer server, DataStreamApi.Resolver readResolver) 
{
     this.server = server;
+    this.readResolver = readResolver;
     this.name = server.getId() + "-" + 
JavaUtils.getClassSimpleName(getClass());
 
     final RaftProperties properties = server.getProperties();
@@ -182,6 +198,22 @@ public class ReadStreamManagement {
       return false;
     }
 
+    if (readResolver != null) {
+      final DataStreamApi dataApi;
+      try {
+        dataApi = readResolver.resolve(request);
+      } catch (Throwable t) {
+        replyDataStreamException(server, t, request, requestBuf, ctx);
+        return true;
+      }
+      if (dataApi != null) {
+        final RaftClientReply reply = 
RaftClientReply.newBuilder().setRequest(request).setSuccess().build();
+        final long streamId = requestBuf.getStreamId();
+        requestExecutor.execute(() -> transfer(dataApi, request, streamId, 
ctx, reply));
+        return true;
+      }
+    }
+
     final RaftServer.Division division;
     try {
       division = server.getDivision(request.getRaftGroupId());
@@ -209,16 +241,32 @@ public class ReadStreamManagement {
         return;
       }
 
-      final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx, readCheckReply);
-      try {
-        division.getStateMachine().data().transferTo(request.getMessage(), 
stream);
-      } catch (Throwable t) {
-        LOG.error("{}: Failed read-only data stream query for {}", this, 
request, t);
-      }
+      transfer(division.getStateMachine().data(), request, 
requestBuf.getStreamId(), ctx, readCheckReply);
     }, requestExecutor);
     return true;
   }
 
+  private void transfer(DataStreamApi dataApi, RaftClientRequest request, long 
streamId,
+      ChannelHandlerContext ctx, RaftClientReply terminalReply) {
+    final ReadStream stream = new ReadStream(request, streamId, ctx, 
terminalReply);
+    try {
+      dataApi.transferTo(request.getMessage(), stream);
+    } catch (Throwable t) {
+      LOG.error("{}: Failed resolved read-only data stream transferTo for {}", 
this, request, t);
+      if (stream.isOpen()) {
+        final RaftClientReply failure = RaftClientReply.newBuilder()
+            .setRequest(request)
+            .setException(new DataStreamException(server.getId(), t))
+            .build();
+        try {
+          stream.fail(failure);
+        } catch (IOException e) {
+          LOG.warn("{}: Failed to send read-only data stream error for {}", 
this, request, e);
+        }
+      }
+    }
+  }
+
   private static RaftClientRequest newDummyReadRequest(RaftClientRequest 
request) {
     final RaftProtos.ReadRequestTypeProto original = 
request.getType().getRead();
     return RaftClientRequest.newBuilder()
diff --git 
a/ratis-server-api/src/main/java/org/apache/ratis/server/RaftServerConfigKeys.java
 
b/ratis-server-api/src/main/java/org/apache/ratis/server/RaftServerConfigKeys.java
index 2d5559478..79e03f554 100644
--- 
a/ratis-server-api/src/main/java/org/apache/ratis/server/RaftServerConfigKeys.java
+++ 
b/ratis-server-api/src/main/java/org/apache/ratis/server/RaftServerConfigKeys.java
@@ -17,7 +17,9 @@
  */
 package org.apache.ratis.server;
 
+import org.apache.ratis.conf.Parameters;
 import org.apache.ratis.conf.RaftProperties;
+import org.apache.ratis.server.api.DataStreamApi;
 import org.apache.ratis.util.JavaUtils;
 import org.apache.ratis.util.SizeInBytes;
 import org.apache.ratis.util.TimeDuration;
@@ -847,6 +849,17 @@ public interface RaftServerConfigKeys {
     static void setClientPoolSize(RaftProperties properties, int num) {
       setInt(properties::setInt, CLIENT_POOL_SIZE_KEY, num);
     }
+
+    String SERVER_API_RESOLVER_PARAMETER = PREFIX + ".server.api.resolver";
+    Class<DataStreamApi.Resolver> SERVER_API_RESOLVER_CLASS = 
DataStreamApi.Resolver.class;
+    static DataStreamApi.Resolver serverApiResolver(Parameters parameters) {
+      return parameters != null
+          ? parameters.get(SERVER_API_RESOLVER_PARAMETER, 
SERVER_API_RESOLVER_CLASS)
+          : null;
+    }
+    static void setServerApiResolver(Parameters parameters, 
DataStreamApi.Resolver resolver) {
+      parameters.put(SERVER_API_RESOLVER_PARAMETER, resolver, 
SERVER_API_RESOLVER_CLASS);
+    }
   }
 
   /** server rpc timeout related */
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
index bf5335063..3c7ef5b26 100644
--- 
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
@@ -18,6 +18,7 @@
 package org.apache.ratis.server.api;
 
 import org.apache.ratis.protocol.Message;
+import org.apache.ratis.protocol.RaftClientRequest;
 
 import java.io.IOException;
 import java.nio.channels.WritableByteChannel;
@@ -34,4 +35,11 @@ public interface DataStreamApi {
    * @return the number of bytes transferred
    */
   long transferTo(Message request, WritableByteChannel stream) throws 
IOException;
+
+  /** For resolving {@link DataStreamApi}. */
+  @FunctionalInterface
+  interface Resolver {
+    /** @return the data API handling this request, or null if the API is not 
found. */
+    DataStreamApi resolve(RaftClientRequest request) throws IOException;
+  }
 }
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 4d311706a..fd7fae6bc 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
@@ -128,6 +128,177 @@ class TestDataStreamManagement {
     }
   }
 
+  @Test
+  void readResolverHandlesRequestWithoutDivision() throws Exception {
+    final RaftPeerId serverId = RaftPeerId.valueOf("s1");
+    final ClientId clientId = ClientId.randomId();
+    final RaftGroupId groupId = RaftGroupId.randomId();
+    final ByteString query = ByteString.copyFromUtf8("query");
+    final ByteString response = ByteString.copyFromUtf8("response");
+    final AtomicBoolean readCheckSubmitted = new AtomicBoolean();
+    final AtomicReference<RaftClientRequest> resolvedRequest = new 
AtomicReference<>();
+    final AtomicReference<Message> messageRef = new AtomicReference<>();
+    final AtomicReference<WritableByteChannel> streamRef = new 
AtomicReference<>();
+
+    final DataApi dataApi = new DataApi() {
+      @Override
+      public long transferTo(Message request, WritableByteChannel stream) {
+        messageRef.set(request);
+        streamRef.set(stream);
+        return 0;
+      }
+    };
+    final RaftServer server = newRaftServer(serverId, new RaftProperties(), 
null, null, request -> {
+      readCheckSubmitted.set(true);
+      return successReply(request);
+    });
+    final ReadStreamManagement management = new ReadStreamManagement(server, 
request -> {
+      resolvedRequest.set(request);
+      return dataApi;
+    });
+    final EmbeddedChannel embeddedChannel = new EmbeddedChannel(new 
ChannelInboundHandlerAdapter());
+    final ReadOnlyRequest readOnlyRequest = newReadOnlyRequest(clientId, 
serverId, groupId, 1L, query);
+
+    try {
+      assertTrue(management.process(readOnlyRequest.request, 
embeddedChannel.pipeline().firstContext()));
+      assertEquals(0, readOnlyRequest.headerBuf.refCnt());
+
+      JavaUtils.attempt(() -> assertNotNull(streamRef.get()), 10,
+          TimeDuration.valueOf(100, TimeUnit.MILLISECONDS), "resolved 
read-only stream", null);
+      streamRef.get().write(response.asReadOnlyByteBuffer());
+      streamRef.get().close();
+
+      final List<DataStreamReply> replies = new ArrayList<>();
+      JavaUtils.attempt(() -> {
+        for (Object outbound; (outbound = embeddedChannel.readOutbound()) != 
null;) {
+          replies.add((DataStreamReply) outbound);
+        }
+        assertEquals(2, replies.size());
+      }, 10, TimeDuration.valueOf(100, TimeUnit.MILLISECONDS), "resolved 
read-only replies", null);
+
+      assertEquals(query, resolvedRequest.get().getMessage().getContent());
+      assertEquals(query, messageRef.get().getContent());
+      assertFalse(readCheckSubmitted.get(), "a resolved request should bypass 
the Raft read check");
+      assertSuccessReply(Type.STREAM_DATA, response.size(), replies.get(0));
+      assertSuccessReply(Type.STREAM_HEADER, 0, replies.get(1));
+      
assertTrue(ClientProtoUtils.getRaftClientReply(replies.get(1)).isSuccess());
+    } finally {
+      embeddedChannel.finishAndReleaseAll();
+      management.shutdown();
+    }
+  }
+
+  @Test
+  void declinedReadResolverUsesDivisionAndReadCheck() throws Exception {
+    final RaftPeerId serverId = RaftPeerId.valueOf("s1");
+    final ClientId clientId = ClientId.randomId();
+    final RaftGroupId groupId = RaftGroupId.randomId();
+    final AtomicBoolean resolverCalled = new AtomicBoolean();
+    final AtomicReference<RaftClientRequest> submittedReadCheck = new 
AtomicReference<>();
+    final AtomicReference<WritableByteChannel> streamRef = new 
AtomicReference<>();
+
+    final DataApi dataApi = new DataApi() {
+      @Override
+      public long transferTo(Message request, WritableByteChannel stream) {
+        streamRef.set(stream);
+        return 0;
+      }
+    };
+    final StateMachine stateMachine = new BaseStateMachine() {
+      @Override
+      public DataApi data() {
+        return dataApi;
+      }
+    };
+    final RaftServer server = newRaftServer(serverId, new RaftProperties(), 
groupId, newDivision(stateMachine),
+        request -> {
+          submittedReadCheck.set(request);
+          return successReply(request);
+        });
+    final ReadStreamManagement management = new ReadStreamManagement(server, 
request -> {
+      resolverCalled.set(true);
+      return null;
+    });
+    final EmbeddedChannel embeddedChannel = new EmbeddedChannel(new 
ChannelInboundHandlerAdapter());
+    final ReadOnlyRequest readOnlyRequest = newReadOnlyRequest(
+        clientId, serverId, groupId, 1L, ByteString.copyFromUtf8("query"));
+
+    try {
+      assertTrue(management.process(readOnlyRequest.request, 
embeddedChannel.pipeline().firstContext()));
+      JavaUtils.attempt(() -> assertNotNull(streamRef.get()), 10,
+          TimeDuration.valueOf(100, TimeUnit.MILLISECONDS), "declined 
read-only stream", null);
+
+      assertTrue(resolverCalled.get());
+      assertEquals(OrderedAsync.DUMMY.getContent(), 
submittedReadCheck.get().getMessage().getContent());
+      streamRef.get().close();
+    } finally {
+      embeddedChannel.finishAndReleaseAll();
+      management.shutdown();
+    }
+  }
+
+  @Test
+  void readResolverExceptionReturnsFailure() throws Exception {
+    final RaftPeerId serverId = RaftPeerId.valueOf("s1");
+    final ClientId clientId = ClientId.randomId();
+    final RaftGroupId groupId = RaftGroupId.randomId();
+    final RaftServer server = newRaftServer(serverId, new RaftProperties());
+    final ReadStreamManagement management = new ReadStreamManagement(server, 
request -> {
+      throw new IOException("resolve failed");
+    });
+    final EmbeddedChannel embeddedChannel = new EmbeddedChannel(new 
ChannelInboundHandlerAdapter());
+    final ReadOnlyRequest readOnlyRequest = newReadOnlyRequest(
+        clientId, serverId, groupId, 1L, ByteString.copyFromUtf8("query"));
+
+    try {
+      assertTrue(management.process(readOnlyRequest.request, 
embeddedChannel.pipeline().firstContext()));
+      final DataStreamReply reply = embeddedChannel.readOutbound();
+      assertNotNull(reply);
+      assertEquals(Type.STREAM_HEADER, reply.getType());
+      assertFalse(reply.isSuccess());
+      
assertNotNull(ClientProtoUtils.getRaftClientReply(reply).getDataStreamException());
+    } finally {
+      embeddedChannel.finishAndReleaseAll();
+      management.shutdown();
+    }
+  }
+
+  @Test
+  void resolvedQueryExceptionReturnsFailure() throws Exception {
+    final RaftPeerId serverId = RaftPeerId.valueOf("s1");
+    final ClientId clientId = ClientId.randomId();
+    final RaftGroupId groupId = RaftGroupId.randomId();
+    final DataApi dataApi = new DataApi() {
+      @Override
+      public long transferTo(Message request, WritableByteChannel stream) {
+        throw new IllegalStateException("query failed");
+      }
+    };
+    final RaftServer server = newRaftServer(serverId, new RaftProperties());
+    final ReadStreamManagement management = new ReadStreamManagement(server, 
request -> dataApi);
+    final EmbeddedChannel embeddedChannel = new EmbeddedChannel(new 
ChannelInboundHandlerAdapter());
+    final ReadOnlyRequest readOnlyRequest = newReadOnlyRequest(
+        clientId, serverId, groupId, 1L, ByteString.copyFromUtf8("query"));
+
+    try {
+      assertTrue(management.process(readOnlyRequest.request, 
embeddedChannel.pipeline().firstContext()));
+      assertEquals(0, readOnlyRequest.headerBuf.refCnt());
+      final AtomicReference<DataStreamReply> replyRef = new 
AtomicReference<>();
+      JavaUtils.attempt(() -> {
+        replyRef.set(embeddedChannel.readOutbound());
+        assertNotNull(replyRef.get());
+      }, 10, TimeDuration.valueOf(100, TimeUnit.MILLISECONDS), "resolved query 
failure", null);
+
+      final DataStreamReply reply = replyRef.get();
+      assertEquals(Type.STREAM_HEADER, reply.getType());
+      assertFalse(reply.isSuccess());
+      
assertNotNull(ClientProtoUtils.getRaftClientReply(reply).getDataStreamException());
+    } finally {
+      embeddedChannel.finishAndReleaseAll();
+      management.shutdown();
+    }
+  }
+
   @Test
   void readOnlyRequestWaitsForLinearizableCheck() throws Exception {
     final RaftPeerId serverId = RaftPeerId.valueOf("s1");

Reply via email to