2022-08-12 14:24:32 +08:00
|
|
|
// Copyright 2022 Memgraph Ltd.
|
|
|
|
//
|
|
|
|
// Use of this software is governed by the Business Source License
|
|
|
|
// included in the file licenses/BSL.txt; by using this file, you agree to be bound by the terms of the Business Source
|
|
|
|
// License, and you may not use this file except in compliance with the Business Source License.
|
|
|
|
//
|
|
|
|
// As of the Change Date specified in that file, in accordance with
|
|
|
|
// the Business Source License, use of this software will be governed
|
|
|
|
// by the Apache License, Version 2.0, included in the file
|
|
|
|
// licenses/APL.txt.
|
|
|
|
|
2022-11-18 01:36:46 +08:00
|
|
|
#include <memory>
|
2022-08-12 14:24:32 +08:00
|
|
|
#include <thread>
|
|
|
|
|
2022-11-18 01:36:46 +08:00
|
|
|
#include <gmock/gmock.h>
|
|
|
|
#include <gtest/gtest.h>
|
|
|
|
#include <spdlog/cfg/env.h>
|
|
|
|
|
|
|
|
#include "io/message_histogram_collector.hpp"
|
2022-08-12 14:24:32 +08:00
|
|
|
#include "io/simulator/simulator.hpp"
|
2022-10-26 21:57:11 +08:00
|
|
|
#include "utils/print_helpers.hpp"
|
2022-08-12 14:24:32 +08:00
|
|
|
|
|
|
|
using memgraph::io::Address;
|
|
|
|
using memgraph::io::Io;
|
2022-11-18 01:36:46 +08:00
|
|
|
using memgraph::io::LatencyHistogramSummaries;
|
2022-08-12 14:24:32 +08:00
|
|
|
using memgraph::io::ResponseFuture;
|
|
|
|
using memgraph::io::ResponseResult;
|
|
|
|
using memgraph::io::simulator::Simulator;
|
|
|
|
using memgraph::io::simulator::SimulatorConfig;
|
2022-11-18 01:36:46 +08:00
|
|
|
using memgraph::io::simulator::SimulatorStats;
|
2022-08-12 14:24:32 +08:00
|
|
|
using memgraph::io::simulator::SimulatorTransport;
|
|
|
|
|
|
|
|
struct CounterRequest {
|
|
|
|
uint64_t proposal;
|
|
|
|
};
|
|
|
|
|
|
|
|
struct CounterResponse {
|
|
|
|
uint64_t highest_seen;
|
|
|
|
};
|
|
|
|
|
|
|
|
void run_server(Io<SimulatorTransport> io) {
|
|
|
|
uint64_t highest_seen = 0;
|
|
|
|
|
|
|
|
while (!io.ShouldShutDown()) {
|
|
|
|
std::cout << "[SERVER] Is receiving..." << std::endl;
|
|
|
|
auto request_result = io.Receive<CounterRequest>();
|
|
|
|
if (request_result.HasError()) {
|
|
|
|
std::cout << "[SERVER] Error, continue" << std::endl;
|
|
|
|
continue;
|
|
|
|
}
|
2022-11-18 01:36:46 +08:00
|
|
|
std::cout << "[SERVER] Got message" << std::endl;
|
2022-08-12 14:24:32 +08:00
|
|
|
auto request_envelope = request_result.GetValue();
|
|
|
|
auto req = std::get<CounterRequest>(request_envelope.message);
|
|
|
|
|
|
|
|
highest_seen = std::max(highest_seen, req.proposal);
|
|
|
|
auto srv_res = CounterResponse{highest_seen};
|
|
|
|
|
|
|
|
io.Send(request_envelope.from_address, request_envelope.request_id, srv_res);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-11-18 01:36:46 +08:00
|
|
|
std::pair<SimulatorStats, LatencyHistogramSummaries> RunWorkload(SimulatorConfig &config) {
|
2022-08-12 14:24:32 +08:00
|
|
|
auto simulator = Simulator(config);
|
|
|
|
|
|
|
|
auto cli_addr = Address::TestAddress(1);
|
|
|
|
auto srv_addr = Address::TestAddress(2);
|
|
|
|
|
|
|
|
Io<SimulatorTransport> cli_io = simulator.Register(cli_addr);
|
|
|
|
Io<SimulatorTransport> srv_io = simulator.Register(srv_addr);
|
|
|
|
|
|
|
|
auto srv_thread = std::jthread(run_server, std::move(srv_io));
|
|
|
|
simulator.IncrementServerCountAndWaitForQuiescentState(srv_addr);
|
|
|
|
|
|
|
|
for (int i = 1; i < 3; ++i) {
|
|
|
|
// send request
|
|
|
|
CounterRequest cli_req;
|
|
|
|
cli_req.proposal = i;
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("[CLIENT] calling Request");
|
2022-08-12 14:24:32 +08:00
|
|
|
auto res_f = cli_io.Request<CounterRequest, CounterResponse>(srv_addr, cli_req);
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("[CLIENT] calling Wait");
|
2022-08-12 14:24:32 +08:00
|
|
|
auto res_rez = std::move(res_f).Wait();
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("[CLIENT] Wait returned");
|
2022-08-12 14:24:32 +08:00
|
|
|
if (!res_rez.HasError()) {
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("[CLIENT] Got a valid response");
|
2022-08-12 14:24:32 +08:00
|
|
|
auto env = res_rez.GetValue();
|
|
|
|
MG_ASSERT(env.message.highest_seen == i);
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("response latency: {} microseconds", env.response_latency.count());
|
2022-08-12 14:24:32 +08:00
|
|
|
} else {
|
2022-11-18 01:36:46 +08:00
|
|
|
spdlog::info("[CLIENT] Got an error");
|
2022-08-12 14:24:32 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
simulator.ShutDown();
|
2022-11-18 01:36:46 +08:00
|
|
|
|
|
|
|
return std::make_pair(simulator.Stats(), cli_io.ResponseLatencies());
|
|
|
|
}
|
|
|
|
|
|
|
|
int main() {
|
|
|
|
spdlog::cfg::load_env_levels();
|
|
|
|
|
|
|
|
auto config = SimulatorConfig{
|
|
|
|
.drop_percent = 0,
|
|
|
|
.perform_timeouts = true,
|
|
|
|
.scramble_messages = true,
|
|
|
|
.rng_seed = 0,
|
|
|
|
};
|
|
|
|
|
|
|
|
auto [sim_stats_1, latency_stats_1] = RunWorkload(config);
|
|
|
|
auto [sim_stats_2, latency_stats_2] = RunWorkload(config);
|
|
|
|
|
|
|
|
if (sim_stats_1 != sim_stats_2 || latency_stats_1 != latency_stats_2) {
|
2022-11-18 05:22:41 +08:00
|
|
|
spdlog::error("simulator stats diverged across runs");
|
|
|
|
spdlog::error("run 1 simulator stats: {}", sim_stats_1);
|
|
|
|
spdlog::error("run 2 simulator stats: {}", sim_stats_2);
|
|
|
|
spdlog::error("run 1 latency:\n{}", latency_stats_1.SummaryTable());
|
|
|
|
spdlog::error("run 2 latency:\n{}", latency_stats_2.SummaryTable());
|
2022-11-18 01:36:46 +08:00
|
|
|
std::terminate();
|
|
|
|
}
|
|
|
|
|
2022-08-12 14:24:32 +08:00
|
|
|
return 0;
|
|
|
|
}
|