Hoist read and write requests on the RSM into wrapper structs for futureproofing
This commit is contained in:
@@ -28,8 +28,10 @@ using memgraph::io::ResponseEnvelope;
|
||||
using memgraph::io::ResponseFuture;
|
||||
using memgraph::io::ResponseResult;
|
||||
using memgraph::io::rsm::Raft;
|
||||
// using memgraph::io::rsm::ReplicationRequest;
|
||||
// using memgraph::io::rsm::ReplicationResponse;
|
||||
using memgraph::io::rsm::ReadRequest;
|
||||
using memgraph::io::rsm::ReadResponse;
|
||||
using memgraph::io::rsm::WriteRequest;
|
||||
using memgraph::io::rsm::WriteResponse;
|
||||
using memgraph::io::simulator::Simulator;
|
||||
using memgraph::io::simulator::SimulatorConfig;
|
||||
using memgraph::io::simulator::SimulatorStats;
|
||||
@@ -44,24 +46,22 @@ struct CasRequest {
|
||||
struct CasResponse {
|
||||
bool success;
|
||||
std::optional<int> last_value;
|
||||
std::optional<Address> retry_leader;
|
||||
};
|
||||
|
||||
struct ReadRequest {
|
||||
struct GetRequest {
|
||||
int key;
|
||||
};
|
||||
|
||||
struct ReadResponse {
|
||||
struct GetResponse {
|
||||
int value;
|
||||
std::optional<Address> retry_leader;
|
||||
};
|
||||
|
||||
class TestState {
|
||||
std::map<int, int> state_;
|
||||
|
||||
public:
|
||||
ReadResponse read(ReadRequest request) {
|
||||
ReadResponse ret;
|
||||
GetResponse read(GetRequest request) {
|
||||
GetResponse ret;
|
||||
ret.value = state_.at(request.key);
|
||||
return ret;
|
||||
}
|
||||
@@ -113,7 +113,7 @@ class TestState {
|
||||
};
|
||||
|
||||
template <typename IoImpl>
|
||||
void RunRaft(Raft<IoImpl, TestState, CasRequest, CasResponse, ReadRequest, ReadResponse> server) {
|
||||
void RunRaft(Raft<IoImpl, TestState, CasRequest, CasResponse, GetRequest, GetResponse> server) {
|
||||
server.Run();
|
||||
}
|
||||
|
||||
@@ -144,7 +144,7 @@ void RunSimulation() {
|
||||
std::vector<Address> srv_3_peers = {srv_addr_1, srv_addr_2};
|
||||
|
||||
// TODO(tyler / gabor) supply default TestState to Raft constructor
|
||||
using RaftClass = Raft<SimulatorTransport, TestState, CasRequest, CasResponse, ReadRequest, ReadResponse>;
|
||||
using RaftClass = Raft<SimulatorTransport, TestState, CasRequest, CasResponse, GetRequest, GetResponse>;
|
||||
RaftClass srv_1{std::move(srv_io_1), srv_1_peers, TestState{}};
|
||||
RaftClass srv_2{std::move(srv_io_2), srv_2_peers, TestState{}};
|
||||
RaftClass srv_3{std::move(srv_io_3), srv_3_peers, TestState{}};
|
||||
@@ -167,27 +167,30 @@ void RunSimulation() {
|
||||
|
||||
while (true) {
|
||||
// send request
|
||||
CasRequest cli_req;
|
||||
cli_req.key = 0;
|
||||
cli_req.old_value = std::nullopt;
|
||||
cli_req.new_value = 12;
|
||||
CasRequest cas_req;
|
||||
cas_req.key = 0;
|
||||
cas_req.old_value = std::nullopt;
|
||||
cas_req.new_value = 12;
|
||||
|
||||
WriteRequest<CasRequest> cli_req;
|
||||
cli_req.operation = cas_req;
|
||||
|
||||
// TODO(tyler / gabor) replace Replication* with Cas/Read
|
||||
|
||||
std::cout << "client sending ReplicationRequest to Leader " << leader.last_known_port << std::endl;
|
||||
ResponseFuture<CasResponse> response_future =
|
||||
cli_io.RequestWithTimeout<CasRequest, CasResponse>(leader, cli_req, 50000);
|
||||
ResponseFuture<WriteResponse<CasResponse>> response_future =
|
||||
cli_io.RequestWithTimeout<WriteRequest<CasRequest>, WriteResponse<CasResponse>>(leader, cli_req, 50000);
|
||||
|
||||
// receive response
|
||||
ResponseResult<CasResponse> response_result = response_future.Wait();
|
||||
ResponseResult<WriteResponse<CasResponse>> response_result = response_future.Wait();
|
||||
|
||||
if (response_result.HasError()) {
|
||||
std::cout << "client timed out while trying to communicate with leader server " << std::endl;
|
||||
continue;
|
||||
}
|
||||
|
||||
ResponseEnvelope<CasResponse> response_envelope = response_result.GetValue();
|
||||
CasResponse response = response_envelope.message;
|
||||
ResponseEnvelope<WriteResponse<CasResponse>> response_envelope = response_result.GetValue();
|
||||
WriteResponse<CasResponse> response = response_envelope.message;
|
||||
|
||||
if (response.success) {
|
||||
success = true;
|
||||
|
||||
Reference in New Issue
Block a user