diff --git a/build.c b/build.c index b59cda4..0016154 100644 --- a/build.c +++ b/build.c @@ -9,7 +9,7 @@ int main(int argc, char **argv) { char *CXX = strcpy(calloc(1024, 1), or_else(getenv("CXX"), "g++")); 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"}; strcat(CXXFLAGS, " -march=native -O3 -Wpedantic -Wall -Wextra -Wsign-conversion -Wconversion -std=c++20 -Isrc -IuSockets/src"); diff --git a/examples/HelloWorldThreaded.cpp b/examples/HelloWorldThreaded.cpp index e4d871f..42137d0 100644 --- a/examples/HelloWorldThreaded.cpp +++ b/examples/HelloWorldThreaded.cpp @@ -1,40 +1,24 @@ #include "App.h" -#include -#include -#include - -/* Note that SSL is disabled unless you build with WITH_OPENSSL=1 */ -const int SSL = 1; -std::mutex stdoutMutex; +#include "LocalCluster.h" int main() { - /* Overly simple hello world app, using multiple threads */ - std::vector threads(std::thread::hardware_concurrency()); - - std::transform(threads.begin(), threads.end(), threads.begin(), [](std::thread */*t*/) { - return new std::thread([]() { - - uWS::SSLApp({ - .key_file_name = "misc/key.pem", - .cert_file_name = "misc/cert.pem", - .passphrase = "1234" - }).get("/*", [](auto *res, auto * /*req*/) { - res->end("Hello world!"); - }).listen(3000, [](auto *listen_socket) { - stdoutMutex.lock(); - 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; - } - stdoutMutex.unlock(); - }).run(); - + /* Note that SSL is disabled unless you build with WITH_OPENSSL=1 */ + 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(true, (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::for_each(threads.begin(), threads.end(), [](std::thread *t) { - t->join(); - }); } diff --git a/examples/LoadBalancer.cpp b/examples/LoadBalancer.cpp deleted file mode 100644 index 8d76891..0000000 --- a/examples/LoadBalancer.cpp +++ /dev/null @@ -1,101 +0,0 @@ -#include "App.h" -#include -#include -#include - -/* 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 threads(hardwareConcurrency); -std::vector apps; -std::mutex m; - -namespace uWS { -struct LocalCluster { - - //std::vector threads = std::thread::hardware_concurrency(); - std::vector apps; - std::mutex m; - - - static void loadBalancer() { - static std::atomic roundRobin = 0; // atomic fetch_add - } - - LocalCluster(SocketContextOptions options = {}, std::function 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(); - }); -} diff --git a/src/LocalCluster.h b/src/LocalCluster.h new file mode 100644 index 0000000..55ec802 --- /dev/null +++ b/src/LocalCluster.h @@ -0,0 +1,62 @@ +/* This header is highly experimental and needs refactorings but will do for now */ + +#include +#include +#include + +unsigned int roundRobin = 0; +unsigned int hardwareConcurrency = std::thread::hardware_concurrency(); +std::vector threads(hardwareConcurrency); +std::vector apps; +std::mutex m; + +namespace uWS { +struct LocalCluster { + + //std::vector threads = std::thread::hardware_concurrency(); + //std::vector apps; + //std::mutex m; + + + static void loadBalancer() { + static std::atomic roundRobin = 0; // atomic fetch_add + } + + LocalCluster(SocketContextOptions options = {}, std::function 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(); + }); + } +}; +} \ No newline at end of file