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");