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]