This is an automated email from the ASF dual-hosted git repository. ascherbakov pushed a commit to branch ignite-14149 in repository https://gitbox.apache.org/repos/asf/ignite-3.git
commit ac96caef28b5f0d21fdf1cc53be5abe89de00bf9 Author: Alexey Scherbakov <[email protected]> AuthorDate: Sat Feb 27 13:52:39 2021 +0300 IGNITE-14149 wip. --- ...CommonMessages.java => RaftClientMessages.java} | 6 +- .../client/message/AddLearnersRequestImpl.java | 6 +- .../raft/client/message/AddPeerRequestImpl.java | 6 +- .../raft/client/message/AddPeerResponseImpl.java | 6 +- .../raft/client/message/ChangePeerRequestImpl.java | 6 +- .../client/message/ChangePeersResponseImpl.java | 6 +- .../message/ClientMessageBuilderFactory.java | 44 +++---- .../raft/client/message/GetLeaderRequestImpl.java | 6 +- .../raft/client/message/GetLeaderResponseImpl.java | 6 +- .../raft/client/message/GetPeersRequestImpl.java | 6 +- .../raft/client/message/GetPeersResponseImpl.java | 6 +- .../client/message/LearnersOpResponseImpl.java | 6 +- .../raft/client/message/PingRequestImpl.java | 6 +- .../RaftClientCommonMessageBuilderFactory.java | 84 ------------- .../message/RaftClientMessageBuilderFactory.java | 94 +++++++++++++++ .../client/message/RemoveLearnersRequestImpl.java | 6 +- .../raft/client/message/RemovePeerRequestImpl.java | 6 +- .../client/message/RemovePeerResponseImpl.java | 6 +- .../client/message/ResetLearnersRequestImpl.java | 6 +- .../raft/client/message/ResetPeerRequestImpl.java | 6 +- .../raft/client/message/SnapshotRequestImpl.java | 7 +- .../raft/client/message/StatusResponseImpl.java | 6 +- .../client/message/TransferLeaderRequestImpl.java | 6 +- .../raft/client/message/UserRequestImpl.java | 32 +++++ .../raft/client/message/UserResponseImpl.java | 21 ++++ .../ignite/raft/client/rpc/RaftGroupRpcClient.java | 34 +++--- .../client/rpc/impl/RaftGroupRpcClientImpl.java | 26 ++-- .../client/service/RaftGroupManagmentService.java | 8 +- .../impl/RaftGroupClientRequestServiceImpl.java | 13 +- .../client/{ => rpc}/RaftGroupRpcClientTest.java | 25 ++-- .../service/RaftGroupClientRequestServiceTest.java | 132 +++++++++++++++++++++ 31 files changed, 416 insertions(+), 218 deletions(-) diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientCommonMessages.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientMessages.java similarity index 98% rename from modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientCommonMessages.java rename to modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientMessages.java index 5a40dc6..3f44685 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientCommonMessages.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/RaftClientMessages.java @@ -27,8 +27,8 @@ import org.apache.ignite.raft.rpc.RaftGroupMessage; /** * */ -public final class RaftClientCommonMessages { - private RaftClientCommonMessages() { +public final class RaftClientMessages { + private RaftClientMessages() { } public interface StatusResponse extends Message { @@ -263,6 +263,8 @@ public final class RaftClientCommonMessages { public interface Builder<T> { Builder setRequest(T request); + Builder setGroupId(String groupId); + UserRequest<T> build(); } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddLearnersRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddLearnersRequestImpl.java index cc3a763..7261c54 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddLearnersRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddLearnersRequestImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -public class AddLearnersRequestImpl implements RaftClientCommonMessages.AddLearnersRequest, RaftClientCommonMessages.AddLearnersRequest.Builder { +public class AddLearnersRequestImpl implements RaftClientMessages.AddLearnersRequest, RaftClientMessages.AddLearnersRequest.Builder { private String groupId; private PeerId leaderId; private List<PeerId> learnersList = new ArrayList<>(); @@ -31,7 +31,7 @@ public class AddLearnersRequestImpl implements RaftClientCommonMessages.AddLearn return this; } - @Override public RaftClientCommonMessages.AddLearnersRequest build() { + @Override public RaftClientMessages.AddLearnersRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerRequestImpl.java index 01fb298..ba6d2c8 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerRequestImpl.java @@ -1,9 +1,9 @@ package org.apache.ignite.raft.client.message; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class AddPeerRequestImpl implements RaftClientCommonMessages.AddPeerRequest, RaftClientCommonMessages.AddPeerRequest.Builder { +class AddPeerRequestImpl implements RaftClientMessages.AddPeerRequest, RaftClientMessages.AddPeerRequest.Builder { private String groupId; private PeerId leaderId; private PeerId peerId; @@ -28,7 +28,7 @@ class AddPeerRequestImpl implements RaftClientCommonMessages.AddPeerRequest, Raf return this; } - @Override public RaftClientCommonMessages.AddPeerRequest build() { + @Override public RaftClientMessages.AddPeerRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerResponseImpl.java index 127124a..b9e1c00 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/AddPeerResponseImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class AddPeerResponseImpl implements RaftClientCommonMessages.AddPeerResponse, RaftClientCommonMessages.AddPeerResponse.Builder { +class AddPeerResponseImpl implements RaftClientMessages.AddPeerResponse, RaftClientMessages.AddPeerResponse.Builder { private List<PeerId> oldPeersList = new ArrayList<>(); private List<PeerId> newPeersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class AddPeerResponseImpl implements RaftClientCommonMessages.AddPeerResponse, R return this; } - @Override public RaftClientCommonMessages.AddPeerResponse build() { + @Override public RaftClientMessages.AddPeerResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeerRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeerRequestImpl.java index 06570ea..47d3abc 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeerRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeerRequestImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class ChangePeerRequestImpl implements RaftClientCommonMessages.ChangePeersRequest, RaftClientCommonMessages.ChangePeersRequest.Builder { +class ChangePeerRequestImpl implements RaftClientMessages.ChangePeersRequest, RaftClientMessages.ChangePeersRequest.Builder { private String groupId; private List<PeerId> newPeersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class ChangePeerRequestImpl implements RaftClientCommonMessages.ChangePeersReque return this; } - @Override public RaftClientCommonMessages.ChangePeersRequest build() { + @Override public RaftClientMessages.ChangePeersRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeersResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeersResponseImpl.java index d22d8f0..bc9f74d 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeersResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ChangePeersResponseImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class ChangePeersResponseImpl implements RaftClientCommonMessages.ChangePeersResponse, RaftClientCommonMessages.ChangePeersResponse.Builder { +class ChangePeersResponseImpl implements RaftClientMessages.ChangePeersResponse, RaftClientMessages.ChangePeersResponse.Builder { private List<PeerId> oldPeersList = new ArrayList<>(); private List<PeerId> newPeersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class ChangePeersResponseImpl implements RaftClientCommonMessages.ChangePeersRes return this; } - @Override public RaftClientCommonMessages.ChangePeersResponse build() { + @Override public RaftClientMessages.ChangePeersResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ClientMessageBuilderFactory.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ClientMessageBuilderFactory.java index b0f48cf..4351a67 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ClientMessageBuilderFactory.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ClientMessageBuilderFactory.java @@ -1,48 +1,48 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; /** */ public interface ClientMessageBuilderFactory { - public static ClientMessageBuilderFactory DEFAULT_MESSAGE_BUILDER_FACTORY = new RaftClientCommonMessageBuilderFactory(); + RaftClientMessages.PingRequest.Builder createPingRequest(); - RaftClientCommonMessages.PingRequest.Builder createPingRequest(); + RaftClientMessages.StatusResponse.Builder createStatusResponse(); - RaftClientCommonMessages.StatusResponse.Builder createStatusResponse(); + RaftClientMessages.AddPeerRequest.Builder createAddPeerRequest(); - RaftClientCommonMessages.AddPeerRequest.Builder createAddPeerRequest(); + RaftClientMessages.AddPeerResponse.Builder createAddPeerResponse(); - RaftClientCommonMessages.AddPeerResponse.Builder createAddPeerResponse(); + RaftClientMessages.RemovePeerRequest.Builder createRemovePeerRequest(); - RaftClientCommonMessages.RemovePeerRequest.Builder createRemovePeerRequest(); + RaftClientMessages.RemovePeerResponse.Builder createRemovePeerResponse(); - RaftClientCommonMessages.RemovePeerResponse.Builder createRemovePeerResponse(); + RaftClientMessages.ChangePeersRequest.Builder createChangePeerRequest(); - RaftClientCommonMessages.ChangePeersRequest.Builder createChangePeerRequest(); + RaftClientMessages.ChangePeersResponse.Builder createChangePeerResponse(); - RaftClientCommonMessages.ChangePeersResponse.Builder createChangePeerResponse(); + RaftClientMessages.SnapshotRequest.Builder createSnapshotRequest(); - RaftClientCommonMessages.SnapshotRequest.Builder createSnapshotRequest(); + RaftClientMessages.ResetPeerRequest.Builder createResetPeerRequest(); - RaftClientCommonMessages.ResetPeerRequest.Builder createResetPeerRequest(); + RaftClientMessages.TransferLeaderRequest.Builder createTransferLeaderRequest(); - RaftClientCommonMessages.TransferLeaderRequest.Builder createTransferLeaderRequest(); + RaftClientMessages.GetLeaderRequest.Builder createGetLeaderRequest(); - RaftClientCommonMessages.GetLeaderRequest.Builder createGetLeaderRequest(); + RaftClientMessages.GetLeaderResponse.Builder createGetLeaderResponse(); - RaftClientCommonMessages.GetLeaderResponse.Builder createGetLeaderResponse(); + RaftClientMessages.GetPeersRequest.Builder createGetPeersRequest(); - RaftClientCommonMessages.GetPeersRequest.Builder createGetPeersRequest(); + RaftClientMessages.GetPeersResponse.Builder createGetPeersResponse(); - RaftClientCommonMessages.GetPeersResponse.Builder createGetPeersResponse(); + RaftClientMessages.AddLearnersRequest.Builder createAddLearnersRequest(); - RaftClientCommonMessages.AddLearnersRequest.Builder createAddLearnersRequest(); + RaftClientMessages.RemoveLearnersRequest.Builder createRemoveLearnersRequest(); - RaftClientCommonMessages.RemoveLearnersRequest.Builder createRemoveLearnersRequest(); + RaftClientMessages.ResetLearnersRequest.Builder createResetLearnersRequest(); - RaftClientCommonMessages.ResetLearnersRequest.Builder createResetLearnersRequest(); + RaftClientMessages.LearnersOpResponse.Builder createLearnersOpResponse(); - RaftClientCommonMessages.LearnersOpResponse.Builder createLearnersOpResponse(); + RaftClientMessages.UserRequest.Builder createUserRequest(); - RaftClientCommonMessages.UserRequest.Builder createUserRequest(); + RaftClientMessages.UserResponse.Builder createUserResponse(); } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderRequestImpl.java index 429bab4..dbc7566 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderRequestImpl.java @@ -1,9 +1,9 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -public class GetLeaderRequestImpl implements RaftClientCommonMessages.GetLeaderRequest, RaftClientCommonMessages.GetLeaderRequest.Builder { +public class GetLeaderRequestImpl implements RaftClientMessages.GetLeaderRequest, RaftClientMessages.GetLeaderRequest.Builder { private String groupId; @Override public String getGroupId() { @@ -16,7 +16,7 @@ public class GetLeaderRequestImpl implements RaftClientCommonMessages.GetLeaderR return this; } - @Override public RaftClientCommonMessages.GetLeaderRequest build() { + @Override public RaftClientMessages.GetLeaderRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderResponseImpl.java index 22789b1..0c84245 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetLeaderResponseImpl.java @@ -2,16 +2,16 @@ package org.apache.ignite.raft.client.message; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -public class GetLeaderResponseImpl implements RaftClientCommonMessages.GetLeaderResponse, RaftClientCommonMessages.GetLeaderResponse.Builder { +public class GetLeaderResponseImpl implements RaftClientMessages.GetLeaderResponse, RaftClientMessages.GetLeaderResponse.Builder { private PeerId leaderId; @Override public PeerId getLeaderId() { return leaderId; } - @Override public RaftClientCommonMessages.GetLeaderResponse build() { + @Override public RaftClientMessages.GetLeaderResponse build() { return this; } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersRequestImpl.java index 1851152..ea50a57 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersRequestImpl.java @@ -1,8 +1,8 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class GetPeersRequestImpl implements RaftClientCommonMessages.GetPeersRequest, RaftClientCommonMessages.GetPeersRequest.Builder { +class GetPeersRequestImpl implements RaftClientMessages.GetPeersRequest, RaftClientMessages.GetPeersRequest.Builder { private String groupId; private boolean onlyAlive; @@ -26,7 +26,7 @@ class GetPeersRequestImpl implements RaftClientCommonMessages.GetPeersRequest, R return this; } - @Override public RaftClientCommonMessages.GetPeersRequest build() { + @Override public RaftClientMessages.GetPeersRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersResponseImpl.java index 37a4912..bed9b6e 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/GetPeersResponseImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class GetPeersResponseImpl implements RaftClientCommonMessages.GetPeersResponse, RaftClientCommonMessages.GetPeersResponse.Builder { +class GetPeersResponseImpl implements RaftClientMessages.GetPeersResponse, RaftClientMessages.GetPeersResponse.Builder { private List<PeerId> peersList = new ArrayList<>(); private List<PeerId> learnersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class GetPeersResponseImpl implements RaftClientCommonMessages.GetPeersResponse, return this; } - @Override public RaftClientCommonMessages.GetPeersResponse build() { + @Override public RaftClientMessages.GetPeersResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/LearnersOpResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/LearnersOpResponseImpl.java index 4c3f4c1..b929e90 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/LearnersOpResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/LearnersOpResponseImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class LearnersOpResponseImpl implements RaftClientCommonMessages.LearnersOpResponse, RaftClientCommonMessages.LearnersOpResponse.Builder { +class LearnersOpResponseImpl implements RaftClientMessages.LearnersOpResponse, RaftClientMessages.LearnersOpResponse.Builder { private List<PeerId> oldLearnersList = new ArrayList<>(); private List<PeerId> newLearnersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class LearnersOpResponseImpl implements RaftClientCommonMessages.LearnersOpRespo return this; } - @Override public RaftClientCommonMessages.LearnersOpResponse build() { + @Override public RaftClientMessages.LearnersOpResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/PingRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/PingRequestImpl.java index d4ad03d..b0c0a78 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/PingRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/PingRequestImpl.java @@ -1,8 +1,8 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class PingRequestImpl implements RaftClientCommonMessages.PingRequest, RaftClientCommonMessages.PingRequest.Builder { +class PingRequestImpl implements RaftClientMessages.PingRequest, RaftClientMessages.PingRequest.Builder { private long sendTimestamp; @Override public long getSendTimestamp() { @@ -15,7 +15,7 @@ class PingRequestImpl implements RaftClientCommonMessages.PingRequest, RaftClien return this; } - @Override public RaftClientCommonMessages.PingRequest build() { + @Override public RaftClientMessages.PingRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientCommonMessageBuilderFactory.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientCommonMessageBuilderFactory.java deleted file mode 100644 index 440cee2..0000000 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientCommonMessageBuilderFactory.java +++ /dev/null @@ -1,84 +0,0 @@ -package org.apache.ignite.raft.client.message; - -import org.apache.ignite.raft.client.RaftClientCommonMessages; - -/** - * Raft client message builders factory. - */ -public class RaftClientCommonMessageBuilderFactory implements ClientMessageBuilderFactory { - @Override public RaftClientCommonMessages.PingRequest.Builder createPingRequest() { - return new PingRequestImpl(); - } - - @Override public RaftClientCommonMessages.StatusResponse.Builder createStatusResponse() { - return new StatusResponseImpl(); - } - - @Override public RaftClientCommonMessages.AddPeerRequest.Builder createAddPeerRequest() { - return new AddPeerRequestImpl(); - } - - @Override public RaftClientCommonMessages.AddPeerResponse.Builder createAddPeerResponse() { - return new AddPeerResponseImpl(); - } - - @Override public RaftClientCommonMessages.RemovePeerRequest.Builder createRemovePeerRequest() { - return new RemovePeerRequestImpl(); - } - - @Override public RaftClientCommonMessages.RemovePeerResponse.Builder createRemovePeerResponse() { - return new RemovePeerResponseImpl(); - } - - @Override public RaftClientCommonMessages.ChangePeersRequest.Builder createChangePeerRequest() { - return new ChangePeerRequestImpl(); - } - - @Override public RaftClientCommonMessages.ChangePeersResponse.Builder createChangePeerResponse() { - return new ChangePeersResponseImpl(); - } - - @Override public RaftClientCommonMessages.SnapshotRequest.Builder createSnapshotRequest() { - return new SnapshotRequestImpl(); - } - - @Override public RaftClientCommonMessages.ResetPeerRequest.Builder createResetPeerRequest() { - return new ResetPeerRequestImpl(); - } - - @Override public RaftClientCommonMessages.TransferLeaderRequest.Builder createTransferLeaderRequest() { - return new TransferLeaderRequestImpl(); - } - - @Override public RaftClientCommonMessages.GetLeaderRequest.Builder createGetLeaderRequest() { - return new GetLeaderRequestImpl(); - } - - @Override public RaftClientCommonMessages.GetLeaderResponse.Builder createGetLeaderResponse() { - return new GetLeaderResponseImpl(); - } - - @Override public RaftClientCommonMessages.GetPeersRequest.Builder createGetPeersRequest() { - return new GetPeersRequestImpl(); - } - - @Override public RaftClientCommonMessages.GetPeersResponse.Builder createGetPeersResponse() { - return new GetPeersResponseImpl(); - } - - @Override public RaftClientCommonMessages.AddLearnersRequest.Builder createAddLearnersRequest() { - return new AddLearnersRequestImpl(); - } - - @Override public RaftClientCommonMessages.RemoveLearnersRequest.Builder createRemoveLearnersRequest() { - return new RemoveLearnersRequestImpl(); - } - - @Override public RaftClientCommonMessages.ResetLearnersRequest.Builder createResetLearnersRequest() { - return new ResetLearnersRequestImpl(); - } - - @Override public RaftClientCommonMessages.LearnersOpResponse.Builder createLearnersOpResponse() { - return new LearnersOpResponseImpl(); - } -} diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientMessageBuilderFactory.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientMessageBuilderFactory.java new file mode 100644 index 0000000..2d47315 --- /dev/null +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RaftClientMessageBuilderFactory.java @@ -0,0 +1,94 @@ +package org.apache.ignite.raft.client.message; + +import org.apache.ignite.raft.client.RaftClientMessages; + +/** + * Raft client message builders factory. + */ +public class RaftClientMessageBuilderFactory implements ClientMessageBuilderFactory { + public static RaftClientMessageBuilderFactory INSTANCE = new RaftClientMessageBuilderFactory(); + + @Override public RaftClientMessages.PingRequest.Builder createPingRequest() { + return new PingRequestImpl(); + } + + @Override public RaftClientMessages.StatusResponse.Builder createStatusResponse() { + return new StatusResponseImpl(); + } + + @Override public RaftClientMessages.AddPeerRequest.Builder createAddPeerRequest() { + return new AddPeerRequestImpl(); + } + + @Override public RaftClientMessages.AddPeerResponse.Builder createAddPeerResponse() { + return new AddPeerResponseImpl(); + } + + @Override public RaftClientMessages.RemovePeerRequest.Builder createRemovePeerRequest() { + return new RemovePeerRequestImpl(); + } + + @Override public RaftClientMessages.RemovePeerResponse.Builder createRemovePeerResponse() { + return new RemovePeerResponseImpl(); + } + + @Override public RaftClientMessages.ChangePeersRequest.Builder createChangePeerRequest() { + return new ChangePeerRequestImpl(); + } + + @Override public RaftClientMessages.ChangePeersResponse.Builder createChangePeerResponse() { + return new ChangePeersResponseImpl(); + } + + @Override public RaftClientMessages.SnapshotRequest.Builder createSnapshotRequest() { + return new SnapshotRequestImpl(); + } + + @Override public RaftClientMessages.ResetPeerRequest.Builder createResetPeerRequest() { + return new ResetPeerRequestImpl(); + } + + @Override public RaftClientMessages.TransferLeaderRequest.Builder createTransferLeaderRequest() { + return new TransferLeaderRequestImpl(); + } + + @Override public RaftClientMessages.GetLeaderRequest.Builder createGetLeaderRequest() { + return new GetLeaderRequestImpl(); + } + + @Override public RaftClientMessages.GetLeaderResponse.Builder createGetLeaderResponse() { + return new GetLeaderResponseImpl(); + } + + @Override public RaftClientMessages.GetPeersRequest.Builder createGetPeersRequest() { + return new GetPeersRequestImpl(); + } + + @Override public RaftClientMessages.GetPeersResponse.Builder createGetPeersResponse() { + return new GetPeersResponseImpl(); + } + + @Override public RaftClientMessages.AddLearnersRequest.Builder createAddLearnersRequest() { + return new AddLearnersRequestImpl(); + } + + @Override public RaftClientMessages.RemoveLearnersRequest.Builder createRemoveLearnersRequest() { + return new RemoveLearnersRequestImpl(); + } + + @Override public RaftClientMessages.ResetLearnersRequest.Builder createResetLearnersRequest() { + return new ResetLearnersRequestImpl(); + } + + @Override public RaftClientMessages.LearnersOpResponse.Builder createLearnersOpResponse() { + return new LearnersOpResponseImpl(); + } + + @Override public RaftClientMessages.UserRequest.Builder createUserRequest() { + return new UserRequestImpl(); + } + + @Override public RaftClientMessages.UserResponse.Builder createUserResponse() { + return new UserResponseImpl(); + } +} diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemoveLearnersRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemoveLearnersRequestImpl.java index de0b056..e55aefc 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemoveLearnersRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemoveLearnersRequestImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class RemoveLearnersRequestImpl implements RaftClientCommonMessages.RemoveLearnersRequest, RaftClientCommonMessages.RemoveLearnersRequest.Builder { +class RemoveLearnersRequestImpl implements RaftClientMessages.RemoveLearnersRequest, RaftClientMessages.RemoveLearnersRequest.Builder { private String groupId; private List<PeerId> learnersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class RemoveLearnersRequestImpl implements RaftClientCommonMessages.RemoveLearne return this; } - @Override public RaftClientCommonMessages.RemoveLearnersRequest build() { + @Override public RaftClientMessages.RemoveLearnersRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerRequestImpl.java index a068f1d..0e4934d 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerRequestImpl.java @@ -1,9 +1,9 @@ package org.apache.ignite.raft.client.message; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class RemovePeerRequestImpl implements RaftClientCommonMessages.RemovePeerRequest, RaftClientCommonMessages.RemovePeerRequest.Builder { +class RemovePeerRequestImpl implements RaftClientMessages.RemovePeerRequest, RaftClientMessages.RemovePeerRequest.Builder { private String groupId; private PeerId peerId; @@ -27,7 +27,7 @@ class RemovePeerRequestImpl implements RaftClientCommonMessages.RemovePeerReques return this; } - @Override public RaftClientCommonMessages.RemovePeerRequest build() { + @Override public RaftClientMessages.RemovePeerRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerResponseImpl.java index 35d9013..00eae18 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/RemovePeerResponseImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class RemovePeerResponseImpl implements RaftClientCommonMessages.RemovePeerResponse, RaftClientCommonMessages.RemovePeerResponse.Builder { +class RemovePeerResponseImpl implements RaftClientMessages.RemovePeerResponse, RaftClientMessages.RemovePeerResponse.Builder { private List<PeerId> oldPeersList = new ArrayList<>(); private List<PeerId> newPeersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class RemovePeerResponseImpl implements RaftClientCommonMessages.RemovePeerRespo return this; } - @Override public RaftClientCommonMessages.RemovePeerResponse build() { + @Override public RaftClientMessages.RemovePeerResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetLearnersRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetLearnersRequestImpl.java index 1e5188d..340a39f 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetLearnersRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetLearnersRequestImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class ResetLearnersRequestImpl implements RaftClientCommonMessages.ResetLearnersRequest, RaftClientCommonMessages.ResetLearnersRequest.Builder { +class ResetLearnersRequestImpl implements RaftClientMessages.ResetLearnersRequest, RaftClientMessages.ResetLearnersRequest.Builder { private String groupId; private List<PeerId> learnersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class ResetLearnersRequestImpl implements RaftClientCommonMessages.ResetLearners return this; } - @Override public RaftClientCommonMessages.ResetLearnersRequest build() { + @Override public RaftClientMessages.ResetLearnersRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetPeerRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetPeerRequestImpl.java index bb9c28b..6c8d4bc 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetPeerRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/ResetPeerRequestImpl.java @@ -3,9 +3,9 @@ package org.apache.ignite.raft.client.message; import java.util.ArrayList; import java.util.List; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class ResetPeerRequestImpl implements RaftClientCommonMessages.ResetPeerRequest, RaftClientCommonMessages.ResetPeerRequest.Builder { +class ResetPeerRequestImpl implements RaftClientMessages.ResetPeerRequest, RaftClientMessages.ResetPeerRequest.Builder { private String groupId; private List<PeerId> newPeersList = new ArrayList<>(); @@ -29,7 +29,7 @@ class ResetPeerRequestImpl implements RaftClientCommonMessages.ResetPeerRequest, return this; } - @Override public RaftClientCommonMessages.ResetPeerRequest build() { + @Override public RaftClientMessages.ResetPeerRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/SnapshotRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/SnapshotRequestImpl.java index 5f1bde1..08a09dc 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/SnapshotRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/SnapshotRequestImpl.java @@ -1,9 +1,8 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class SnapshotRequestImpl implements RaftClientCommonMessages.SnapshotRequest, RaftClientCommonMessages.SnapshotRequest.Builder { +class SnapshotRequestImpl implements RaftClientMessages.SnapshotRequest, RaftClientMessages.SnapshotRequest.Builder { private String groupId; @Override public String getGroupId() { @@ -16,7 +15,7 @@ class SnapshotRequestImpl implements RaftClientCommonMessages.SnapshotRequest, R return this; } - @Override public RaftClientCommonMessages.SnapshotRequest build() { + @Override public RaftClientMessages.SnapshotRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/StatusResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/StatusResponseImpl.java index 5afcf22..6ae0b6b 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/StatusResponseImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/StatusResponseImpl.java @@ -1,8 +1,8 @@ package org.apache.ignite.raft.client.message; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class StatusResponseImpl implements RaftClientCommonMessages.StatusResponse, RaftClientCommonMessages.StatusResponse.Builder { +class StatusResponseImpl implements RaftClientMessages.StatusResponse, RaftClientMessages.StatusResponse.Builder { private int errorCode; private String errorMsg = ""; @@ -26,7 +26,7 @@ class StatusResponseImpl implements RaftClientCommonMessages.StatusResponse, Raf return this; } - @Override public RaftClientCommonMessages.StatusResponse build() { + @Override public RaftClientMessages.StatusResponse build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/TransferLeaderRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/TransferLeaderRequestImpl.java index b125aec..2cc6592 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/TransferLeaderRequestImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/TransferLeaderRequestImpl.java @@ -1,9 +1,9 @@ package org.apache.ignite.raft.client.message; import org.apache.ignite.raft.PeerId; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; -class TransferLeaderRequestImpl implements RaftClientCommonMessages.TransferLeaderRequest, RaftClientCommonMessages.TransferLeaderRequest.Builder { +class TransferLeaderRequestImpl implements RaftClientMessages.TransferLeaderRequest, RaftClientMessages.TransferLeaderRequest.Builder { private String groupId; private PeerId peerId; @@ -21,7 +21,7 @@ class TransferLeaderRequestImpl implements RaftClientCommonMessages.TransferLead return this; } - @Override public RaftClientCommonMessages.TransferLeaderRequest build() { + @Override public RaftClientMessages.TransferLeaderRequest build() { return this; } } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserRequestImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserRequestImpl.java new file mode 100644 index 0000000..06d667e --- /dev/null +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserRequestImpl.java @@ -0,0 +1,32 @@ +package org.apache.ignite.raft.client.message; + +import org.apache.ignite.raft.client.RaftClientMessages; + +public class UserRequestImpl<T> implements RaftClientMessages.UserRequest<T>, RaftClientMessages.UserRequest.Builder<T> { + private T request; + private String groupId; + + @Override public T request() { + return request; + } + + @Override public Builder setGroupId(String groupId) { + this.groupId = groupId; + + return this; + } + + @Override public Builder setRequest(T request) { + this.request = request; + + return this; + } + + @Override public RaftClientMessages.UserRequest<T> build() { + return this; + } + + @Override public String getGroupId() { + return groupId; + } +} diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserResponseImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserResponseImpl.java new file mode 100644 index 0000000..0d62953 --- /dev/null +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/message/UserResponseImpl.java @@ -0,0 +1,21 @@ +package org.apache.ignite.raft.client.message; + +import org.apache.ignite.raft.client.RaftClientMessages; + +public class UserResponseImpl<T> implements RaftClientMessages.UserResponse<T>, RaftClientMessages.UserResponse.Builder<T> { + private T response; + + @Override public T response() { + return response; + } + + @Override public Builder setResponse(T response) { + this.response = response; + + return this; + } + + @Override public RaftClientMessages.UserResponse build() { + return this; + } +} diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClient.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClient.java index 2eca6e1..14384e9 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClient.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClient.java @@ -22,22 +22,21 @@ import org.apache.ignite.raft.PeerId; import org.apache.ignite.raft.client.message.ClientMessageBuilderFactory; import org.apache.ignite.raft.rpc.Message; import org.apache.ignite.raft.rpc.RaftGroupMessage; -import org.jetbrains.annotations.Nullable; - -import static org.apache.ignite.raft.client.RaftClientCommonMessages.AddLearnersRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.AddPeerRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.AddPeerResponse; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.ChangePeersRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.ChangePeersResponse; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.LearnersOpResponse; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.RemoveLearnersRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.RemovePeerRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.RemovePeerResponse; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.ResetLearnersRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.ResetPeerRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.SnapshotRequest; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.StatusResponse; -import static org.apache.ignite.raft.client.RaftClientCommonMessages.TransferLeaderRequest; + +import static org.apache.ignite.raft.client.RaftClientMessages.AddLearnersRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.AddPeerRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.AddPeerResponse; +import static org.apache.ignite.raft.client.RaftClientMessages.ChangePeersRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.ChangePeersResponse; +import static org.apache.ignite.raft.client.RaftClientMessages.LearnersOpResponse; +import static org.apache.ignite.raft.client.RaftClientMessages.RemoveLearnersRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.RemovePeerRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.RemovePeerResponse; +import static org.apache.ignite.raft.client.RaftClientMessages.ResetLearnersRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.ResetPeerRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.SnapshotRequest; +import static org.apache.ignite.raft.client.RaftClientMessages.StatusResponse; +import static org.apache.ignite.raft.client.RaftClientMessages.TransferLeaderRequest; /** * Low-level raft group RPC client. @@ -161,5 +160,8 @@ public interface RaftGroupRpcClient { */ <R extends Message> CompletableFuture<R> sendCustom(RaftGroupMessage request); + /** + * @return A message builder factory. + */ ClientMessageBuilderFactory factory(); } diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/impl/RaftGroupRpcClientImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/impl/RaftGroupRpcClientImpl.java index e9c67cb..4599556 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/impl/RaftGroupRpcClientImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/rpc/impl/RaftGroupRpcClientImpl.java @@ -14,9 +14,9 @@ import java.util.function.BiConsumer; import org.apache.ignite.raft.PeerId; import org.apache.ignite.raft.RaftException; import org.apache.ignite.raft.State; -import org.apache.ignite.raft.client.RaftClientCommonMessages; -import org.apache.ignite.raft.client.RaftClientCommonMessages.GetLeaderResponse; -import org.apache.ignite.raft.client.RaftClientCommonMessages.StatusResponse; +import org.apache.ignite.raft.client.RaftClientMessages; +import org.apache.ignite.raft.client.RaftClientMessages.GetLeaderResponse; +import org.apache.ignite.raft.client.RaftClientMessages.StatusResponse; import org.apache.ignite.raft.client.message.ClientMessageBuilderFactory; import org.apache.ignite.raft.client.rpc.RaftGroupRpcClient; import org.apache.ignite.raft.rpc.InvokeCallback; @@ -87,39 +87,39 @@ public class RaftGroupRpcClientImpl implements RaftGroupRpcClient { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.AddPeerResponse> addPeer(RaftClientCommonMessages.AddPeerRequest request) { + @Override public CompletableFuture<RaftClientMessages.AddPeerResponse> addPeer(RaftClientMessages.AddPeerRequest request) { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.RemovePeerResponse> removePeer(RaftClientCommonMessages.RemovePeerRequest request) { + @Override public CompletableFuture<RaftClientMessages.RemovePeerResponse> removePeer(RaftClientMessages.RemovePeerRequest request) { return null; } - @Override public CompletableFuture<StatusResponse> resetPeers(PeerId peerId, RaftClientCommonMessages.ResetPeerRequest request) { + @Override public CompletableFuture<StatusResponse> resetPeers(PeerId peerId, RaftClientMessages.ResetPeerRequest request) { return null; } - @Override public CompletableFuture<StatusResponse> snapshot(PeerId peerId, RaftClientCommonMessages.SnapshotRequest request) { + @Override public CompletableFuture<StatusResponse> snapshot(PeerId peerId, RaftClientMessages.SnapshotRequest request) { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.ChangePeersResponse> changePeers(RaftClientCommonMessages.ChangePeersRequest request) { + @Override public CompletableFuture<RaftClientMessages.ChangePeersResponse> changePeers(RaftClientMessages.ChangePeersRequest request) { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.LearnersOpResponse> addLearners(RaftClientCommonMessages.AddLearnersRequest request) { + @Override public CompletableFuture<RaftClientMessages.LearnersOpResponse> addLearners(RaftClientMessages.AddLearnersRequest request) { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.LearnersOpResponse> removeLearners(RaftClientCommonMessages.RemoveLearnersRequest request) { + @Override public CompletableFuture<RaftClientMessages.LearnersOpResponse> removeLearners(RaftClientMessages.RemoveLearnersRequest request) { return null; } - @Override public CompletableFuture<RaftClientCommonMessages.LearnersOpResponse> resetLearners(RaftClientCommonMessages.ResetLearnersRequest request) { + @Override public CompletableFuture<RaftClientMessages.LearnersOpResponse> resetLearners(RaftClientMessages.ResetLearnersRequest request) { return null; } - @Override public CompletableFuture<StatusResponse> transferLeader(RaftClientCommonMessages.TransferLeaderRequest request) { + @Override public CompletableFuture<StatusResponse> transferLeader(RaftClientMessages.TransferLeaderRequest request) { return null; } @@ -133,7 +133,7 @@ public class RaftGroupRpcClientImpl implements RaftGroupRpcClient { return fut; if (state.updateFutRef.compareAndSet(null, (fut = new CompletableFuture<>()))) { - RaftClientCommonMessages.GetLeaderRequest req = factory.createGetLeaderRequest().setGroupId(groupId).build(); + RaftClientMessages.GetLeaderRequest req = factory.createGetLeaderRequest().setGroupId(groupId).build(); CompletableFuture<GetLeaderResponse> finalFut = fut; diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/RaftGroupManagmentService.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/RaftGroupManagmentService.java index 725d3e1..d868dc9 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/RaftGroupManagmentService.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/RaftGroupManagmentService.java @@ -9,24 +9,24 @@ import org.jetbrains.annotations.Nullable; public interface RaftGroupManagmentService { /** * @param groupId - * @return Peer id. + * @return Leader id or null if it has not been yet initialized. */ @Nullable PeerId getLeader(String groupId); /** * @param groupId - * @return List of peers. + * @return List of peers or null if it has not been yet initialized. */ @Nullable List<PeerId> getPeers(String groupId); /** * @param groupId - * @return List of peers. + * @return List of peers or null if it has not been yet initialized. */ @Nullable List<PeerId> getLearners(String groupId); /** - * Adds a voring peer to the raft group. + * Adds a voting peer to the raft group. * * @param request request data * @return A future with the result diff --git a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/impl/RaftGroupClientRequestServiceImpl.java b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/impl/RaftGroupClientRequestServiceImpl.java index b5439c2..16d63c4 100644 --- a/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/impl/RaftGroupClientRequestServiceImpl.java +++ b/modules/raft-client/src/main/java/org/apache/ignite/raft/client/service/impl/RaftGroupClientRequestServiceImpl.java @@ -2,24 +2,27 @@ package org.apache.ignite.raft.client.service.impl; import java.util.concurrent.CompletableFuture; import java.util.function.Function; -import org.apache.ignite.raft.client.RaftClientCommonMessages; +import org.apache.ignite.raft.client.RaftClientMessages; import org.apache.ignite.raft.client.rpc.RaftGroupRpcClient; import org.apache.ignite.raft.client.service.RaftGroupClientRequestService; import org.apache.ignite.raft.rpc.Message; public class RaftGroupClientRequestServiceImpl implements RaftGroupClientRequestService { - private RaftGroupRpcClient rpcClient; + private final RaftGroupRpcClient rpcClient; + private final String groupId; - public RaftGroupClientRequestServiceImpl(RaftGroupRpcClient rpcClient) { + public RaftGroupClientRequestServiceImpl(RaftGroupRpcClient rpcClient, String groupId) { this.rpcClient = rpcClient; + this.groupId = groupId; } @Override public <T, R> CompletableFuture<R> submit(T request) { - RaftClientCommonMessages.UserRequest r = rpcClient.factory().createUserRequest().setRequest(request).build(); + RaftClientMessages.UserRequest r = + rpcClient.factory().createUserRequest().setRequest(request).setGroupId(groupId).build(); return rpcClient.sendCustom(r).thenApply(new Function<Message, R>() { @Override public R apply(Message message) { - RaftClientCommonMessages.UserResponse<R> resp = (RaftClientCommonMessages.UserResponse<R>) message; + RaftClientMessages.UserResponse<R> resp = (RaftClientMessages.UserResponse<R>) message; return resp.response(); } diff --git a/modules/raft-client/src/test/java/org/apache/ignite/raft/client/RaftGroupRpcClientTest.java b/modules/raft-client/src/test/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClientTest.java similarity index 86% rename from modules/raft-client/src/test/java/org/apache/ignite/raft/client/RaftGroupRpcClientTest.java rename to modules/raft-client/src/test/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClientTest.java index c476422..5247c59 100644 --- a/modules/raft-client/src/test/java/org/apache/ignite/raft/client/RaftGroupRpcClientTest.java +++ b/modules/raft-client/src/test/java/org/apache/ignite/raft/client/rpc/RaftGroupRpcClientTest.java @@ -1,15 +1,13 @@ -package org.apache.ignite.raft.client; +package org.apache.ignite.raft.client.rpc; -import java.util.Collections; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.Executor; import java.util.concurrent.TimeoutException; import org.apache.ignite.raft.PeerId; import org.apache.ignite.raft.State; -import org.apache.ignite.raft.client.RaftClientCommonMessages.GetLeaderRequest; +import org.apache.ignite.raft.client.RaftClientMessages.GetLeaderRequest; import org.apache.ignite.raft.client.rpc.impl.RaftGroupRpcClientImpl; -import org.apache.ignite.raft.client.rpc.RaftGroupRpcClient; import org.apache.ignite.raft.rpc.InvokeCallback; import org.apache.ignite.raft.rpc.Message; import org.apache.ignite.raft.rpc.NodeImpl; @@ -23,7 +21,8 @@ import org.mockito.invocation.InvocationOnMock; import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.stubbing.Answer; -import static org.apache.ignite.raft.client.message.ClientMessageBuilderFactory.DEFAULT_MESSAGE_BUILDER_FACTORY; +import static java.util.Collections.singleton; +import static org.apache.ignite.raft.client.message.RaftClientMessageBuilderFactory.INSTANCE; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -44,8 +43,7 @@ public class RaftGroupRpcClientTest { mockLeaderRequest(false); - RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, DEFAULT_MESSAGE_BUILDER_FACTORY, - 5_000, Collections.singleton(leader.getNode())); + RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, INSTANCE, 5_000, singleton(leader.getNode())); PeerId leaderId = client.refreshLeader(groupId).get(); @@ -59,8 +57,7 @@ public class RaftGroupRpcClientTest { mockLeaderRequest(false); - RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, DEFAULT_MESSAGE_BUILDER_FACTORY, - 5_000, Collections.singleton(leader.getNode())); + RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, INSTANCE, 5_000, singleton(leader.getNode())); int cnt = 20; @@ -104,8 +101,8 @@ public class RaftGroupRpcClientTest { mockLeaderRequest(true); - RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, DEFAULT_MESSAGE_BUILDER_FACTORY, - 5_000, Collections.singleton(leader.getNode())); + RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, INSTANCE, + 5_000, singleton(leader.getNode())); try { client.refreshLeader(groupId).get(); @@ -124,8 +121,8 @@ public class RaftGroupRpcClientTest { mockLeaderRequest(false); mockCustomRequest(); - RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, DEFAULT_MESSAGE_BUILDER_FACTORY, - 5_000, Collections.singleton(leader.getNode())); + RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, INSTANCE, + 5_000, singleton(leader.getNode())); JunkRequest req = new JunkRequest(groupId); @@ -176,7 +173,7 @@ public class RaftGroupRpcClientTest { if (timeout) callback.complete(null, new TimeoutException()); else - callback.complete(DEFAULT_MESSAGE_BUILDER_FACTORY.createGetLeaderResponse().setLeaderId(leader).build(), null); + callback.complete(INSTANCE.createGetLeaderResponse().setLeaderId(leader).build(), null); }); return null; diff --git a/modules/raft-client/src/test/java/org/apache/ignite/raft/client/service/RaftGroupClientRequestServiceTest.java b/modules/raft-client/src/test/java/org/apache/ignite/raft/client/service/RaftGroupClientRequestServiceTest.java new file mode 100644 index 0000000..56e3ff8 --- /dev/null +++ b/modules/raft-client/src/test/java/org/apache/ignite/raft/client/service/RaftGroupClientRequestServiceTest.java @@ -0,0 +1,132 @@ +package org.apache.ignite.raft.client.service; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executor; +import org.apache.ignite.raft.PeerId; +import org.apache.ignite.raft.client.RaftClientMessages; +import org.apache.ignite.raft.client.RaftClientMessages.UserRequest; +import org.apache.ignite.raft.client.rpc.RaftGroupRpcClient; +import org.apache.ignite.raft.client.rpc.impl.RaftGroupRpcClientImpl; +import org.apache.ignite.raft.client.service.impl.RaftGroupClientRequestServiceImpl; +import org.apache.ignite.raft.rpc.InvokeCallback; +import org.apache.ignite.raft.rpc.NodeImpl; +import org.apache.ignite.raft.rpc.RpcClient; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentMatcher; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.stubbing.Answer; + +import static java.util.Collections.singleton; +import static org.apache.ignite.raft.client.message.RaftClientMessageBuilderFactory.INSTANCE; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.ArgumentMatchers.eq; + +@ExtendWith(MockitoExtension.class) +public class RaftGroupClientRequestServiceTest { + @Mock + private RpcClient rpcClient; + + private static PeerId leader = new PeerId(new NodeImpl("test")); + + @Test + public void testUserRequest() throws Exception { + String groupId = "test"; + + mockLeaderRequest(); + mockUserRequest1(); + mockUserRequest2(); + + RaftGroupRpcClient client = new RaftGroupRpcClientImpl(rpcClient, INSTANCE, 5_000, singleton(leader.getNode())); + + RaftGroupClientRequestService service = new RaftGroupClientRequestServiceImpl(client, groupId); + + assertNull(client.state(groupId).leader()); + + CompletableFuture<TestOutput1> fut1 = service.submit(new TestInput1()); + + TestOutput1 output1 = fut1.get(); + + assertNotNull(output1); + + CompletableFuture<TestOutput2> fut2 = service.submit(new TestInput2()); + + TestOutput2 output2 = fut2.get(); + + assertNotNull(output2); + + assertEquals(leader, client.state(groupId).leader()); + } + + private static class TestInput1 { + } + + private static class TestOutput1 { + } + + private static class TestInput2 { + } + + private static class TestOutput2 { + } + + private void mockLeaderRequest() { + Mockito.doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { + InvokeCallback callback = invocation.getArgument(2); + Executor executor = invocation.getArgument(3); + + executor.execute(() -> { + callback.complete(INSTANCE.createGetLeaderResponse().setLeaderId(leader).build(), null); + }); + + return null; + } + }).when(rpcClient).invokeAsync(eq(leader.getNode()), any(RaftClientMessages.GetLeaderRequest.class), + any(), any(), anyLong()); + } + + private void mockUserRequest1() { + Mockito.doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { + InvokeCallback callback = invocation.getArgument(2); + Executor executor = invocation.getArgument(3); + + executor.execute(() -> callback.complete(INSTANCE.createUserResponse(). + setResponse(new TestOutput1()).build(), null)); + + return null; + } + }).when(rpcClient).invokeAsync(eq(leader.getNode()), argThat(new ArgumentMatcher<UserRequest>() { + @Override public boolean matches(UserRequest argument) { + return argument.request() instanceof TestInput1; + } + }), any(), any(), anyLong()); + } + + private void mockUserRequest2() { + Mockito.doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { + InvokeCallback callback = invocation.getArgument(2); + Executor executor = invocation.getArgument(3); + + executor.execute(() -> callback.complete(INSTANCE.createUserResponse(). + setResponse(new TestOutput2()).build(), null)); + + return null; + } + }).when(rpcClient).invokeAsync(eq(leader.getNode()), argThat(new ArgumentMatcher<UserRequest>() { + @Override public boolean matches(UserRequest argument) { + return argument.request() instanceof TestInput2; + } + }), any(), any(), anyLong()); + } +}
