szetszwo commented on code in PR #1490:
URL: https://github.com/apache/ratis/pull/1490#discussion_r3483515802


##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java:
##########
@@ -186,17 +176,80 @@ private boolean processImpl(DataStreamRequestByteBuf 
requestBuf, ChannelHandlerC
       return true;
     }
 
-    final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx);
-    requestExecutor.execute(() -> {
+    final CompletableFuture<RaftClientReply> readOnlyCheck;
+    try {
+      readOnlyCheck = 
server.submitClientRequestAsync(newDummyReadRequest(request));
+    } catch (IOException e) {
+      replyDataStreamException(server, e, request, requestBuf, ctx);
+      return true;
+    }
+
+    readOnlyCheck.whenCompleteAsync((reply, exception) -> {
+      if (exception != null) {
+        replyDataStreamException(server, exception, request, requestBuf, ctx);
+        return;
+      }
+
+      final RaftClientReply terminalReply = toReadStreamReply(request, reply);
+      if (!reply.isSuccess()) {
+        ctx.writeAndFlush(newDataStreamReplyByteBuffer(requestBuf, 
terminalReply));
+        return;
+      }
+
+      final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx, terminalReply);
       try {
         division.getStateMachine().data().query(request.getMessage(), stream);
       } catch (Throwable t) {
         LOG.error("{}: Failed read-only data stream query for {}", this, 
request, t);
       }
-    });
+    }, requestExecutor);
     return true;
   }
 
+  private static RaftClientRequest newDummyReadRequest(RaftClientRequest 
request) {
+    final RaftClientRequest.Builder builder = RaftClientRequest.newBuilder()
+        .setClientId(request.getClientId())
+        .setGroupId(request.getRaftGroupId())
+        .setCallId(request.getCallId())
+        .setMessage(OrderedAsync.DUMMY)
+        .setType(request.getType())
+        .setRepliedCallIds(request.getRepliedCallIds())
+        .setSlidingWindowEntry(request.getSlidingWindowEntry())
+        .setRoutingTable(request.getRoutingTable())
+        .setTimeoutMs(request.getTimeoutMs())
+        .setSpanContext(request.getSpanContext());
+    if (request.isToLeader()) {
+      builder.setLeaderId(request.getServerId());
+    } else {
+      builder.setServerId(request.getServerId());
+    }
+    return builder.build();

Review Comment:
   Let's add a set(RaftClientRequest request) to Builder to copy everything.



##########
ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java:
##########
@@ -1113,11 +1113,12 @@ private CompletableFuture<RaftClientReply> 
readAsync(RaftClientRequest request)
     if (request.getType().getRead().getPreferNonLinearizable()
         || readOption == RaftServerConfigKeys.Read.Option.DEFAULT) {
       final CompletableFuture<RaftClientReply> reply = 
checkLeaderState(request);
-       if (reply != null) {
-         return reply;
-       }
-       return queryStateMachine(request);
-    } else if (readOption == RaftServerConfigKeys.Read.Option.LINEARIZABLE){
+      if (reply != null) {
+        return reply;
+      }
+      return isDummyRead(request) ? 
CompletableFuture.completedFuture(newSuccessReply(request))
+          : queryStateMachine(request);

Review Comment:
   Let's add a DUMMY_SUCCESS_REPLY and move it to queryStateMachine:
   ```java
     static final CompletableFuture<RaftClientReply> DUMMY_SUCCESS_REPLY
         = 
CompletableFuture.completedFuture(RaftClientReply.newBuilder().setSuccess().build());
   ```
   ```java
     CompletableFuture<RaftClientReply> queryStateMachine(RaftClientRequest 
request) {
       if (request.getType().getRead().getDummy()) {
         return DUMMY_SUCCESS_REPLY;
       }
       return processQueryFuture(stateMachine.query(request.getMessage()), 
request);
     }
   ```



##########
ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java:
##########
@@ -1136,12 +1137,17 @@ private CompletableFuture<RaftClientReply> 
readAsync(RaftClientRequest request)
       return replyFuture
           .thenCompose(readIndex -> 
getState().getReadRequests().waitToAdvance(readIndex,
               () -> getReadException("add", 
snapshotInstallationHandler.getInProgressInstallSnapshotIndex(), false)))
-          .thenCompose(readIndex -> queryStateMachine(request))
+          .thenCompose(readIndex -> isDummyRead(request)
+              ? CompletableFuture.completedFuture(newSuccessReply(request)) : 
queryStateMachine(request))
           .exceptionally(e -> readException2Reply(request, e));
     } else {
       throw new IllegalStateException("Unexpected read option: " + readOption);
     }
   }
+  private static boolean isDummyRead(RaftClientRequest request) {
+    return request.getMessage() != null && 
OrderedAsync.DUMMY.getContent().equals(request.getMessage().getContent());
+  }

Review Comment:
   We cannot use DUMMY since user applications can set any message.  They may 
already have used it in their messages.  We probably need to add a flag to the 
proto:
   ```diff
   +++ b/ratis-proto/src/main/proto/Raft.proto
   @@ -302,6 +302,7 @@ message ForwardRequestTypeProto {
    message ReadRequestTypeProto {
      bool preferNonLinearizable = 1;
      bool readAfterWriteConsistent = 2;
   +  bool dummy = 3;
    }
   ```



##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/ReadStreamManagement.java:
##########
@@ -186,17 +176,80 @@ private boolean processImpl(DataStreamRequestByteBuf 
requestBuf, ChannelHandlerC
       return true;
     }
 
-    final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx);
-    requestExecutor.execute(() -> {
+    final CompletableFuture<RaftClientReply> readOnlyCheck;
+    try {
+      readOnlyCheck = 
server.submitClientRequestAsync(newDummyReadRequest(request));
+    } catch (IOException e) {
+      replyDataStreamException(server, e, request, requestBuf, ctx);
+      return true;
+    }
+
+    readOnlyCheck.whenCompleteAsync((reply, exception) -> {
+      if (exception != null) {
+        replyDataStreamException(server, exception, request, requestBuf, ctx);
+        return;
+      }
+
+      final RaftClientReply terminalReply = toReadStreamReply(request, reply);
+      if (!reply.isSuccess()) {
+        ctx.writeAndFlush(newDataStreamReplyByteBuffer(requestBuf, 
terminalReply));
+        return;
+      }

Review Comment:
   Since server returns a DUMMY_SUCCESS_REPLY, we could build the failed reply 
directly:
   ```java
         if (!readCheckReply.isSuccess()) {
           final RaftClientReply reply = RaftClientReply.newBuilder()
               .setRequest(request)
               .setSuccess(false)
               .build();
           ctx.writeAndFlush(newDataStreamReplyByteBuffer(requestBuf, reply));
           return;
         }
   
         final ReadStream stream = new ReadStream(request, 
requestBuf.getStreamId(), ctx);
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to