2017-08-17 20:46:34 +08:00
|
|
|
#include <fstream>
|
2017-10-25 20:47:46 +08:00
|
|
|
#include <iostream>
|
2017-08-17 20:46:34 +08:00
|
|
|
|
2017-08-22 22:29:23 +08:00
|
|
|
#include <glog/logging.h>
|
2017-08-17 20:46:34 +08:00
|
|
|
|
2017-10-13 21:30:06 +08:00
|
|
|
#include "memgraph_config.hpp"
|
|
|
|
#include "reactors_distributed.hpp"
|
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
DEFINE_int64(my_mnid, 0, "Memgraph node id"); // TODO(zuza): this should be
|
|
|
|
// assigned by the leader once in
|
|
|
|
// the future
|
2017-08-17 20:46:34 +08:00
|
|
|
|
2017-08-23 19:48:25 +08:00
|
|
|
class MemgraphDistributed {
|
2017-08-17 20:46:34 +08:00
|
|
|
private:
|
|
|
|
using Location = std::pair<std::string, uint16_t>;
|
|
|
|
|
|
|
|
public:
|
2017-08-22 23:13:12 +08:00
|
|
|
/**
|
|
|
|
* Get the (singleton) instance of MemgraphDistributed.
|
|
|
|
*
|
2017-10-25 20:47:46 +08:00
|
|
|
* More info:
|
|
|
|
* https://stackoverflow.com/questions/1008019/c-singleton-design-pattern
|
2017-08-22 23:13:12 +08:00
|
|
|
*/
|
|
|
|
static MemgraphDistributed &GetInstance() {
|
2017-10-25 20:47:46 +08:00
|
|
|
static MemgraphDistributed
|
|
|
|
memgraph; // guaranteed to be destroyed, initialized on first use
|
2017-08-23 19:48:25 +08:00
|
|
|
return memgraph;
|
2017-08-22 23:13:12 +08:00
|
|
|
}
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
/** Register memgraph node id to the given location. */
|
2017-10-25 20:47:46 +08:00
|
|
|
void RegisterMemgraphNode(int64_t mnid, const std::string &address,
|
|
|
|
uint16_t port) {
|
|
|
|
std::unique_lock<std::mutex> lock(mutex_);
|
2017-08-17 20:46:34 +08:00
|
|
|
mnodes_[mnid] = Location(address, port);
|
|
|
|
}
|
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
EventStream *FindChannel(int64_t mnid, const std::string &reactor,
|
2017-08-17 20:46:34 +08:00
|
|
|
const std::string &channel) {
|
2017-10-25 20:47:46 +08:00
|
|
|
std::unique_lock<std::mutex> lock(mutex_);
|
2017-08-22 22:29:23 +08:00
|
|
|
const auto &location = mnodes_.at(mnid);
|
2017-10-25 20:47:46 +08:00
|
|
|
return Distributed::GetInstance().FindChannel(
|
|
|
|
location.first, location.second, reactor, channel);
|
2017-08-17 20:46:34 +08:00
|
|
|
}
|
|
|
|
|
2017-08-22 23:13:12 +08:00
|
|
|
protected:
|
|
|
|
MemgraphDistributed() {}
|
|
|
|
|
2017-08-17 20:46:34 +08:00
|
|
|
private:
|
2017-10-25 20:47:46 +08:00
|
|
|
std::mutex mutex_;
|
2017-08-17 20:46:34 +08:00
|
|
|
std::unordered_map<int64_t, Location> mnodes_;
|
2017-08-22 23:13:12 +08:00
|
|
|
|
|
|
|
MemgraphDistributed(const MemgraphDistributed &) = delete;
|
|
|
|
MemgraphDistributed(MemgraphDistributed &&) = delete;
|
|
|
|
MemgraphDistributed &operator=(const MemgraphDistributed &) = delete;
|
|
|
|
MemgraphDistributed &operator=(MemgraphDistributed &&) = delete;
|
2017-08-17 20:46:34 +08:00
|
|
|
};
|
|
|
|
|
|
|
|
/**
|
|
|
|
* About config file
|
|
|
|
*
|
|
|
|
* Each line contains three strings:
|
|
|
|
* memgraph node id, ip address of the worker, and port of the worker
|
|
|
|
* Data on the first line is used to start master.
|
|
|
|
* Data on the remaining lines is used to start workers.
|
|
|
|
*/
|
|
|
|
|
|
|
|
/**
|
|
|
|
* Parse config file and register processes into system.
|
|
|
|
*
|
|
|
|
* @return Pair (master mnid, list of worker's id).
|
|
|
|
*/
|
2017-10-25 20:47:46 +08:00
|
|
|
std::pair<int64_t, std::vector<int64_t>> ParseConfigAndRegister(
|
|
|
|
const std::string &filename) {
|
2017-08-17 20:46:34 +08:00
|
|
|
std::ifstream file(filename, std::ifstream::in);
|
|
|
|
assert(file.good());
|
|
|
|
int64_t master_mnid;
|
|
|
|
std::vector<int64_t> worker_mnids;
|
|
|
|
int64_t mnid;
|
2017-08-22 22:29:23 +08:00
|
|
|
std::string address;
|
|
|
|
uint16_t port;
|
|
|
|
file >> master_mnid >> address >> port;
|
2017-08-23 19:48:25 +08:00
|
|
|
MemgraphDistributed &memgraph = MemgraphDistributed::GetInstance();
|
|
|
|
memgraph.RegisterMemgraphNode(master_mnid, address, port);
|
2017-08-17 20:46:34 +08:00
|
|
|
while (file.good()) {
|
2017-08-22 22:29:23 +08:00
|
|
|
file >> mnid >> address >> port;
|
2017-10-25 20:47:46 +08:00
|
|
|
if (file.eof()) break;
|
2017-08-23 19:48:25 +08:00
|
|
|
memgraph.RegisterMemgraphNode(mnid, address, port);
|
2017-08-17 20:46:34 +08:00
|
|
|
worker_mnids.push_back(mnid);
|
2017-08-22 22:29:23 +08:00
|
|
|
}
|
|
|
|
file.close();
|
|
|
|
return std::make_pair(master_mnid, worker_mnids);
|
2017-08-17 20:46:34 +08:00
|
|
|
}
|
|
|
|
|
2017-08-22 22:29:23 +08:00
|
|
|
/**
|
|
|
|
* Sends a text message and has a return address.
|
|
|
|
*/
|
2017-08-23 21:16:26 +08:00
|
|
|
class TextMessage : public ReturnAddressMsg {
|
2017-10-25 20:47:46 +08:00
|
|
|
public:
|
2017-08-22 22:29:23 +08:00
|
|
|
TextMessage(std::string reactor, std::string channel, std::string s)
|
2017-10-25 20:47:46 +08:00
|
|
|
: ReturnAddressMsg(reactor, channel), text(s) {}
|
2017-08-22 22:29:23 +08:00
|
|
|
|
|
|
|
template <class Archive>
|
|
|
|
void serialize(Archive &archive) {
|
2017-08-23 21:16:26 +08:00
|
|
|
archive(cereal::virtual_base_class<ReturnAddressMsg>(this), text);
|
2017-08-22 22:29:23 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
std::string text;
|
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
protected:
|
2017-08-22 22:29:23 +08:00
|
|
|
friend class cereal::access;
|
2017-10-25 20:47:46 +08:00
|
|
|
TextMessage() {} // Cereal needs access to a default constructor.
|
2017-08-22 22:29:23 +08:00
|
|
|
};
|
|
|
|
CEREAL_REGISTER_TYPE(TextMessage);
|
|
|
|
|
2017-08-22 23:13:12 +08:00
|
|
|
class Master : public Reactor {
|
2017-08-17 20:46:34 +08:00
|
|
|
public:
|
2017-08-22 23:13:12 +08:00
|
|
|
Master(std::string name, int64_t mnid, std::vector<int64_t> &&worker_mnids)
|
2017-10-25 20:47:46 +08:00
|
|
|
: Reactor(name), mnid_(mnid), worker_mnids_(std::move(worker_mnids)) {}
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
virtual void Run() {
|
2017-08-23 19:48:25 +08:00
|
|
|
MemgraphDistributed &memgraph = MemgraphDistributed::GetInstance();
|
|
|
|
Distributed &distributed = Distributed::GetInstance();
|
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
std::cout << "Master (" << mnid_ << ") @ "
|
|
|
|
<< distributed.network().Address() << ":"
|
|
|
|
<< distributed.network().Port() << std::endl;
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
auto stream = main_.first;
|
2017-08-22 22:29:23 +08:00
|
|
|
|
2017-08-23 21:16:26 +08:00
|
|
|
// wait until every worker sends a ReturnAddressMsg back, then close
|
2017-10-25 20:47:46 +08:00
|
|
|
stream->OnEvent<TextMessage>(
|
|
|
|
[this](const TextMessage &msg, const Subscription &subscription) {
|
|
|
|
std::cout << "Message from " << msg.Address() << ":" << msg.Port()
|
|
|
|
<< " .. " << msg.text << "\n";
|
|
|
|
++workers_seen;
|
|
|
|
if (workers_seen == static_cast<int64_t>(worker_mnids_.size())) {
|
|
|
|
subscription.Unsubscribe();
|
|
|
|
// Sleep for a while so we can read output in the terminal.
|
|
|
|
// (start_distributed.py runs each process in a new tab which is
|
|
|
|
// closed immediately after process has finished)
|
|
|
|
std::this_thread::sleep_for(std::chrono::seconds(4));
|
|
|
|
CloseChannel("main");
|
|
|
|
}
|
|
|
|
});
|
2017-08-17 20:46:34 +08:00
|
|
|
|
2017-08-22 22:29:23 +08:00
|
|
|
// send a TextMessage to each worker
|
2017-08-17 20:46:34 +08:00
|
|
|
for (auto wmnid : worker_mnids_) {
|
2017-08-23 19:48:25 +08:00
|
|
|
auto stream = memgraph.FindChannel(wmnid, "worker", "main");
|
2017-10-25 20:47:46 +08:00
|
|
|
stream->OnEventOnce().ChainOnce<ChannelResolvedMessage>([this, stream](
|
|
|
|
const ChannelResolvedMessage &msg, const Subscription &) {
|
|
|
|
msg.channelWriter()->Send<TextMessage>("master", "main",
|
|
|
|
"hi from master");
|
|
|
|
stream->Close();
|
|
|
|
});
|
2017-08-17 20:46:34 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
protected:
|
|
|
|
int64_t workers_seen = 0;
|
|
|
|
const int64_t mnid_;
|
|
|
|
std::vector<int64_t> worker_mnids_;
|
|
|
|
};
|
|
|
|
|
2017-08-22 23:13:12 +08:00
|
|
|
class Worker : public Reactor {
|
2017-08-17 20:46:34 +08:00
|
|
|
public:
|
2017-08-22 23:13:12 +08:00
|
|
|
Worker(std::string name, int64_t mnid, int64_t master_mnid)
|
2017-10-25 20:47:46 +08:00
|
|
|
: Reactor(name), mnid_(mnid), master_mnid_(master_mnid) {}
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
virtual void Run() {
|
2017-08-23 19:48:25 +08:00
|
|
|
Distributed &distributed = Distributed::GetInstance();
|
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
std::cout << "Worker (" << mnid_ << ") @ "
|
|
|
|
<< distributed.network().Address() << ":"
|
|
|
|
<< distributed.network().Port() << std::endl;
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
auto stream = main_.first;
|
2017-08-22 22:29:23 +08:00
|
|
|
// wait until master sends us a TextMessage, then reply back and close
|
2017-10-25 20:47:46 +08:00
|
|
|
stream->OnEventOnce().ChainOnce<TextMessage>(
|
|
|
|
[this](const TextMessage &msg, const Subscription &) {
|
|
|
|
std::cout << "Message from " << msg.Address() << ":" << msg.Port()
|
|
|
|
<< " .. " << msg.text << "\n";
|
2017-08-22 22:29:23 +08:00
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
msg.GetReturnChannelWriter()->Send<TextMessage>("worker", "main",
|
|
|
|
"hi from worker");
|
2017-08-22 22:29:23 +08:00
|
|
|
|
2017-10-25 20:47:46 +08:00
|
|
|
// Sleep for a while so we can read output in the terminal.
|
|
|
|
std::this_thread::sleep_for(std::chrono::seconds(4));
|
|
|
|
CloseChannel("main");
|
|
|
|
});
|
2017-08-17 20:46:34 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
protected:
|
|
|
|
const int64_t mnid_;
|
|
|
|
const int64_t master_mnid_;
|
|
|
|
};
|
|
|
|
|
|
|
|
int main(int argc, char *argv[]) {
|
2017-08-22 22:29:23 +08:00
|
|
|
google::InitGoogleLogging(argv[0]);
|
2017-08-17 20:46:34 +08:00
|
|
|
gflags::ParseCommandLineFlags(&argc, &argv, true);
|
|
|
|
|
2017-08-22 22:29:23 +08:00
|
|
|
System &system = System::GetInstance();
|
2017-08-22 23:13:12 +08:00
|
|
|
auto mnids = ParseConfigAndRegister(FLAGS_config_filename);
|
2017-08-23 19:48:25 +08:00
|
|
|
Distributed::GetInstance().StartServices();
|
2017-08-17 20:46:34 +08:00
|
|
|
if (FLAGS_my_mnid == mnids.first)
|
2017-08-22 23:13:12 +08:00
|
|
|
system.Spawn<Master>("master", FLAGS_my_mnid, std::move(mnids.second));
|
2017-08-17 20:46:34 +08:00
|
|
|
else
|
2017-08-22 23:13:12 +08:00
|
|
|
system.Spawn<Worker>("worker", FLAGS_my_mnid, mnids.first);
|
2017-08-17 20:46:34 +08:00
|
|
|
system.AwaitShutdown();
|
2017-08-23 19:48:25 +08:00
|
|
|
Distributed::GetInstance().StopServices();
|
2017-08-17 20:46:34 +08:00
|
|
|
|
|
|
|
return 0;
|
2017-08-17 21:32:12 +08:00
|
|
|
}
|