szetszwo commented on code in PR #1469:
URL: https://github.com/apache/ratis/pull/1469#discussion_r3325852400
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -352,9 +356,36 @@ static DataStreamReplyByteBuffer
newDataStreamReplyByteBuffer(DataStreamRequestB
.setDataStreamPacket(request)
.setBuffer(buffer)
.setSuccess(reply.isSuccess())
+ .setCommitInfos(reply.getCommitInfos())
.build();
}
+ static DataStreamReplyByteBuffer
newDataStreamReadOnlyReplyByteBuffer(DataStreamRequestByteBuf request,
+ long streamOffset, ByteBuffer buffer) {
+ final ByteBuffer readOnlyBuffer = buffer.asReadOnlyBuffer();
+ return DataStreamReplyByteBuffer.newBuilder()
+ .setClientId(request.getClientId())
+ .setType(Type.STREAM_DATA)
+ .setStreamId(request.getStreamId())
+ .setStreamOffset(streamOffset)
+ .setBuffer(readOnlyBuffer)
+ .setSuccess(true)
+ .setBytesWritten(readOnlyBuffer.remaining())
+ .build();
+ }
+
+ private static CompletableFuture<Void> writeAndFlush(ChannelHandlerContext
ctx, DataStreamReply reply) {
+ final CompletableFuture<Void> future = new CompletableFuture<>();
+ ctx.writeAndFlush(reply).addListener(channelFuture -> {
+ if (channelFuture.isSuccess()) {
+ future.complete(null);
+ } else {
+ future.completeExceptionally(channelFuture.cause());
+ }
+ });
+ return future;
+ }
Review Comment:
Let's move all the new code to a new class, say ReadStreamManagement.
##########
dev-support/checkstyle.xml:
##########
@@ -60,6 +60,10 @@
</module>
<module name="SuppressWarningsFilter"/>
+ <module name="SuppressionSingleFilter">
+ <property name="checks" value="FileLength"/>
+ <property name="files" value="RaftServerImpl.java"/>
+ </module>
Review Comment:
Since this PR no longer changes RaftServerImpl. Please revert the
checkstyle file.
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -450,6 +481,23 @@ private void readImpl(DataStreamRequestByteBuf request,
ChannelHandlerContext ct
// add to ChannelMap
channels.add(channelId, key);
+ if (request.getType() == Type.STREAM_HEADER) {
+ final RaftClientRequest raftClientRequest = toRaftClientRequest(request);
+ if (raftClientRequest.is(TypeCase.READ)) {
+ submitReadOnlyRequest(request, raftClientRequest,
ctx).whenComplete((v, exception) -> {
+ try {
+ if (exception != null) {
+ replyDataStreamException(server, exception, raftClientRequest,
request, ctx);
+ }
+ } finally {
+ request.release();
+ channels.remove(channelId, key);
Review Comment:
- request can be release earlier -- we only need to clientId and streamId in
the read stream.
- channelId is not used. So, we don't need a ChannelMap in
ReadStreamManagement.
##########
ratis-netty/src/main/java/org/apache/ratis/netty/server/DataStreamManagement.java:
##########
@@ -450,6 +481,23 @@ private void readImpl(DataStreamRequestByteBuf request,
ChannelHandlerContext ct
// add to ChannelMap
channels.add(channelId, key);
+ if (request.getType() == Type.STREAM_HEADER) {
+ final RaftClientRequest raftClientRequest = toRaftClientRequest(request);
+ if (raftClientRequest.is(TypeCase.READ)) {
+ submitReadOnlyRequest(request, raftClientRequest,
ctx).whenComplete((v, exception) -> {
+ try {
+ if (exception != null) {
+ replyDataStreamException(server, exception, raftClientRequest,
request, ctx);
+ }
+ } finally {
+ request.release();
+ channels.remove(channelId, key);
+ }
+ });
+ return;
+ }
+ }
Review Comment:
The new read streams have nothing to do with the existing write streams.
Let's do the check in NettyServerStreamRpc.
```java
@@ -235,6 +237,9 @@ public class NettyServerStreamRpc implements
DataStreamServerRpc {
final DataStreamRequestByteBuf request =
(DataStreamRequestByteBuf)msg;
try(UncheckedAutoCloseable autoReset = requestRef.set(request)) {
+ if (reads.process(request, ctx)) {
+ return;
+ }
requests.read(request, ctx,
proxies.get(request)::getDataStreamOutput);
}
}
```
--
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]