Expose new LocalCluster interface and use in example
This commit is contained in:
@@ -9,7 +9,7 @@ int main(int argc, char **argv) {
|
|||||||
char *CXX = strcpy(calloc(1024, 1), or_else(getenv("CXX"), "g++"));
|
char *CXX = strcpy(calloc(1024, 1), or_else(getenv("CXX"), "g++"));
|
||||||
char *EXEC_SUFFIX = strcpy(calloc(1024, 1), maybe(getenv("EXEC_SUFFIX")));
|
char *EXEC_SUFFIX = strcpy(calloc(1024, 1), maybe(getenv("EXEC_SUFFIX")));
|
||||||
|
|
||||||
char *EXAMPLE_FILES[] = {"LoadBalancer", "Http3Server", "Broadcast", "HelloWorld", "Crc32", "ServerName",
|
char *EXAMPLE_FILES[] = {"HelloWorldThreaded", "Http3Server", "Broadcast", "HelloWorld", "Crc32", "ServerName",
|
||||||
"EchoServer", "BroadcastingEchoServer", "UpgradeSync", "UpgradeAsync", "ParameterRoutes"};
|
"EchoServer", "BroadcastingEchoServer", "UpgradeSync", "UpgradeAsync", "ParameterRoutes"};
|
||||||
|
|
||||||
strcat(CXXFLAGS, " -march=native -O3 -Wpedantic -Wall -Wextra -Wsign-conversion -Wconversion -std=c++20 -Isrc -IuSockets/src");
|
strcat(CXXFLAGS, " -march=native -O3 -Wpedantic -Wall -Wextra -Wsign-conversion -Wconversion -std=c++20 -Isrc -IuSockets/src");
|
||||||
|
|||||||
@@ -1,40 +1,24 @@
|
|||||||
#include "App.h"
|
#include "App.h"
|
||||||
#include <thread>
|
#include "LocalCluster.h"
|
||||||
#include <algorithm>
|
|
||||||
#include <mutex>
|
|
||||||
|
|
||||||
/* Note that SSL is disabled unless you build with WITH_OPENSSL=1 */
|
|
||||||
const int SSL = 1;
|
|
||||||
std::mutex stdoutMutex;
|
|
||||||
|
|
||||||
int main() {
|
int main() {
|
||||||
/* Overly simple hello world app, using multiple threads */
|
/* Note that SSL is disabled unless you build with WITH_OPENSSL=1 */
|
||||||
std::vector<std::thread *> threads(std::thread::hardware_concurrency());
|
uWS::LocalCluster({
|
||||||
|
.key_file_name = "misc/key.pem",
|
||||||
std::transform(threads.begin(), threads.end(), threads.begin(), [](std::thread */*t*/) {
|
.cert_file_name = "misc/cert.pem",
|
||||||
return new std::thread([]() {
|
.passphrase = "1234"
|
||||||
|
},
|
||||||
uWS::SSLApp({
|
[](uWS::SSLApp &app) {
|
||||||
.key_file_name = "misc/key.pem",
|
/* Here this App instance is defined */
|
||||||
.cert_file_name = "misc/cert.pem",
|
app.get("/*", [](auto *res, auto * /*req*/) {
|
||||||
.passphrase = "1234"
|
res->end("Hello world!");
|
||||||
}).get("/*", [](auto *res, auto * /*req*/) {
|
}).listen(3000, [](auto *listen_socket) {
|
||||||
res->end("Hello world!");
|
if (listen_socket) {
|
||||||
}).listen(3000, [](auto *listen_socket) {
|
/* Note that us_listen_socket_t is castable to us_socket_t */
|
||||||
stdoutMutex.lock();
|
std::cout << "Thread " << std::this_thread::get_id() << " listening on port " << us_socket_local_port(true, (struct us_socket_t *) listen_socket) << std::endl;
|
||||||
if (listen_socket) {
|
} else {
|
||||||
/* Note that us_listen_socket_t is castable to us_socket_t */
|
std::cout << "Thread " << std::this_thread::get_id() << " failed to listen on port 3000" << std::endl;
|
||||||
std::cout << "Thread " << std::this_thread::get_id() << " listening on port " << us_socket_local_port(SSL, (struct us_socket_t *) listen_socket) << std::endl;
|
}
|
||||||
} else {
|
|
||||||
std::cout << "Thread " << std::this_thread::get_id() << " failed to listen on port 3000" << std::endl;
|
|
||||||
}
|
|
||||||
stdoutMutex.unlock();
|
|
||||||
}).run();
|
|
||||||
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
std::for_each(threads.begin(), threads.end(), [](std::thread *t) {
|
|
||||||
t->join();
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,101 +0,0 @@
|
|||||||
#include "App.h"
|
|
||||||
#include <thread>
|
|
||||||
#include <algorithm>
|
|
||||||
#include <mutex>
|
|
||||||
|
|
||||||
/* Note that SSL is disabled unless you build with WITH_OPENSSL=1 */
|
|
||||||
const int SSL = 1;
|
|
||||||
|
|
||||||
unsigned int roundRobin = 0;
|
|
||||||
unsigned int hardwareConcurrency = std::thread::hardware_concurrency();
|
|
||||||
std::vector<std::thread *> threads(hardwareConcurrency);
|
|
||||||
std::vector<uWS::SSLApp *> apps;
|
|
||||||
std::mutex m;
|
|
||||||
|
|
||||||
namespace uWS {
|
|
||||||
struct LocalCluster {
|
|
||||||
|
|
||||||
//std::vector<std::thread *> threads = std::thread::hardware_concurrency();
|
|
||||||
std::vector<uWS::SSLApp *> apps;
|
|
||||||
std::mutex m;
|
|
||||||
|
|
||||||
|
|
||||||
static void loadBalancer() {
|
|
||||||
static std::atomic<unsigned int> roundRobin = 0; // atomic fetch_add
|
|
||||||
}
|
|
||||||
|
|
||||||
LocalCluster(SocketContextOptions options = {}, std::function<void(uWS::SSLApp &)> cb = nullptr) {
|
|
||||||
|
|
||||||
}
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
int main() {
|
|
||||||
|
|
||||||
// can be strictly round robin or not
|
|
||||||
|
|
||||||
// uWS::LocalCluster({
|
|
||||||
// .key_file_name = "misc/key.pem",
|
|
||||||
// .cert_file_name = "misc/cert.pem",
|
|
||||||
// .passphrase = "1234"
|
|
||||||
// },
|
|
||||||
// [](uWS::SSLApp &app) {
|
|
||||||
// /* Here this App instance is defined */
|
|
||||||
// app.get("/*", [](auto *res, auto * /*req*/) {
|
|
||||||
// res->end("Hello world!");
|
|
||||||
// }).listen(3000, [](auto *listen_socket) {
|
|
||||||
// if (listen_socket) {
|
|
||||||
// /* Note that us_listen_socket_t is castable to us_socket_t */
|
|
||||||
// std::cout << "Thread " << std::this_thread::get_id() << " listening on port " << us_socket_local_port(SSL, (struct us_socket_t *) listen_socket) << std::endl;
|
|
||||||
// } else {
|
|
||||||
// std::cout << "Thread " << std::this_thread::get_id() << " failed to listen on port 3000" << std::endl;
|
|
||||||
// }
|
|
||||||
// });
|
|
||||||
// });
|
|
||||||
|
|
||||||
std::transform(threads.begin(), threads.end(), threads.begin(), [](std::thread *) {
|
|
||||||
|
|
||||||
return new std::thread([]() {
|
|
||||||
|
|
||||||
// lock this
|
|
||||||
m.lock();
|
|
||||||
apps.emplace_back(new uWS::SSLApp({
|
|
||||||
.key_file_name = "misc/key.pem",
|
|
||||||
.cert_file_name = "misc/cert.pem",
|
|
||||||
.passphrase = "1234"
|
|
||||||
}));
|
|
||||||
uWS::SSLApp *app = apps.back();
|
|
||||||
|
|
||||||
app->get("/*", [](auto *res, auto * /*req*/) {
|
|
||||||
res->end("Hello world!");
|
|
||||||
}).listen(3000, [](auto *listen_socket) {
|
|
||||||
if (listen_socket) {
|
|
||||||
/* Note that us_listen_socket_t is castable to us_socket_t */
|
|
||||||
std::cout << "Thread " << std::this_thread::get_id() << " listening on port " << us_socket_local_port(SSL, (struct us_socket_t *) listen_socket) << std::endl;
|
|
||||||
} else {
|
|
||||||
std::cout << "Thread " << std::this_thread::get_id() << " failed to listen on port 3000" << std::endl;
|
|
||||||
}
|
|
||||||
}).preOpen([](LIBUS_SOCKET_DESCRIPTOR fd) {
|
|
||||||
|
|
||||||
/* Distribute this socket in round robin fashion */
|
|
||||||
std::cout << "About to load balance " << fd << " to " << roundRobin << std::endl;
|
|
||||||
|
|
||||||
auto receivingApp = apps[roundRobin];
|
|
||||||
apps[roundRobin]->getLoop()->defer([fd, receivingApp]() {
|
|
||||||
receivingApp->adoptSocket(fd);
|
|
||||||
});
|
|
||||||
|
|
||||||
roundRobin = (roundRobin + 1) % hardwareConcurrency;
|
|
||||||
return -1;
|
|
||||||
});
|
|
||||||
m.unlock();
|
|
||||||
app->run();
|
|
||||||
std::cout << "Fallthrough!" << std::endl;
|
|
||||||
delete app;
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
std::for_each(threads.begin(), threads.end(), [](std::thread *t) {
|
|
||||||
t->join();
|
|
||||||
});
|
|
||||||
}
|
|
||||||
@@ -0,0 +1,62 @@
|
|||||||
|
/* This header is highly experimental and needs refactorings but will do for now */
|
||||||
|
|
||||||
|
#include <thread>
|
||||||
|
#include <algorithm>
|
||||||
|
#include <mutex>
|
||||||
|
|
||||||
|
unsigned int roundRobin = 0;
|
||||||
|
unsigned int hardwareConcurrency = std::thread::hardware_concurrency();
|
||||||
|
std::vector<std::thread *> threads(hardwareConcurrency);
|
||||||
|
std::vector<uWS::SSLApp *> apps;
|
||||||
|
std::mutex m;
|
||||||
|
|
||||||
|
namespace uWS {
|
||||||
|
struct LocalCluster {
|
||||||
|
|
||||||
|
//std::vector<std::thread *> threads = std::thread::hardware_concurrency();
|
||||||
|
//std::vector<uWS::SSLApp *> apps;
|
||||||
|
//std::mutex m;
|
||||||
|
|
||||||
|
|
||||||
|
static void loadBalancer() {
|
||||||
|
static std::atomic<unsigned int> roundRobin = 0; // atomic fetch_add
|
||||||
|
}
|
||||||
|
|
||||||
|
LocalCluster(SocketContextOptions options = {}, std::function<void(uWS::SSLApp &)> cb = nullptr) {
|
||||||
|
std::transform(threads.begin(), threads.end(), threads.begin(), [options, &cb](std::thread *) {
|
||||||
|
|
||||||
|
return new std::thread([options, &cb]() {
|
||||||
|
|
||||||
|
// lock this
|
||||||
|
m.lock();
|
||||||
|
apps.emplace_back(new uWS::SSLApp(options));
|
||||||
|
uWS::SSLApp *app = apps.back();
|
||||||
|
|
||||||
|
cb(*app);
|
||||||
|
|
||||||
|
app->preOpen([](LIBUS_SOCKET_DESCRIPTOR fd) {
|
||||||
|
|
||||||
|
/* Distribute this socket in round robin fashion */
|
||||||
|
//std::cout << "About to load balance " << fd << " to " << roundRobin << std::endl;
|
||||||
|
|
||||||
|
auto receivingApp = apps[roundRobin];
|
||||||
|
apps[roundRobin]->getLoop()->defer([fd, receivingApp]() {
|
||||||
|
receivingApp->adoptSocket(fd);
|
||||||
|
});
|
||||||
|
|
||||||
|
roundRobin = (roundRobin + 1) % hardwareConcurrency;
|
||||||
|
return -1;
|
||||||
|
});
|
||||||
|
m.unlock();
|
||||||
|
app->run();
|
||||||
|
std::cout << "Fallthrough!" << std::endl;
|
||||||
|
delete app;
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
std::for_each(threads.begin(), threads.end(), [](std::thread *t) {
|
||||||
|
t->join();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user