From 8e00882bb3df7684c31af57a33c2e6bb8fdb1bf7 Mon Sep 17 00:00:00 2001 From: IlushaShurupov Date: Sun, 16 Jul 2023 20:03:14 +0300 Subject: [PATCH] Networking Initial --- .gitmodules | 3 + Externals/asio | 1 + Storage/CMakeLists.txt | 12 +-- Storage/applications/CMakeLists.txt | 11 +++ Storage/applications/Client.cpp | 96 +++++++++++++------ Storage/applications/Server.cpp | 141 ++++++++++++++++++---------- 6 files changed, 178 insertions(+), 86 deletions(-) create mode 100644 .gitmodules create mode 160000 Externals/asio create mode 100644 Storage/applications/CMakeLists.txt diff --git a/.gitmodules b/.gitmodules new file mode 100644 index 0000000..3f056fd --- /dev/null +++ b/.gitmodules @@ -0,0 +1,3 @@ +[submodule "Externals/asio"] + path = Externals/asio + url = https://github.com/IlyaShurupov/asio.git diff --git a/Externals/asio b/Externals/asio new file mode 160000 index 0000000..d6b95c0 --- /dev/null +++ b/Externals/asio @@ -0,0 +1 @@ +Subproject commit d6b95c0188e0359a8cdbdb6571f0cbacf11a538c diff --git a/Storage/CMakeLists.txt b/Storage/CMakeLists.txt index 28e593f..d18d496 100644 --- a/Storage/CMakeLists.txt +++ b/Storage/CMakeLists.txt @@ -12,12 +12,6 @@ add_library(${PROJECT_NAME} STATIC ${SOURCES} ${HEADERS}) target_include_directories(${PROJECT_NAME} PUBLIC ./public/) target_link_libraries(${PROJECT_NAME} PUBLIC Strings) -### ---------------------- Applications --------------------- ### -add_executable(Server ./applications/Server.cpp) -add_executable(Client ./applications/Client.cpp) -target_link_libraries(Server PUBLIC ${PROJECT_NAME}) -target_link_libraries(Client PUBLIC ${PROJECT_NAME}) - ### -------------------------- Tests -------------------------- ### enable_testing() file(GLOB TEST_SOURCES "./tests/*.cpp") @@ -25,4 +19,8 @@ add_executable(${PROJECT_NAME}Tests ${TEST_SOURCES}) target_link_libraries(${PROJECT_NAME}Tests ${PROJECT_NAME} Utils) add_test(NAME ${PROJECT_NAME}Tests COMMAND ${PROJECT_NAME}Tests) -install(TARGETS ${PROJECT_NAME} LIBRARY DESTINATION ${CMAKE_INSTALL_PREFIX}/${PROJECT_NAME}/lib) \ No newline at end of file +install(TARGETS ${PROJECT_NAME} LIBRARY DESTINATION ${CMAKE_INSTALL_PREFIX}/${PROJECT_NAME}/lib) + + +# todo :remove +add_subdirectory(applications) \ No newline at end of file diff --git a/Storage/applications/CMakeLists.txt b/Storage/applications/CMakeLists.txt new file mode 100644 index 0000000..8e1e2df --- /dev/null +++ b/Storage/applications/CMakeLists.txt @@ -0,0 +1,11 @@ + +cmake_minimum_required(VERSION 3.2) + +project(Applications) + +### ---------------------- Applications --------------------- ### +add_executable(Server ./Server.cpp) +add_executable(Client ./Client.cpp) + +include_directories(Client ./../../Externals/asio/asio/include) +include_directories(Server ./../../Externals/asio/asio/include) \ No newline at end of file diff --git a/Storage/applications/Client.cpp b/Storage/applications/Client.cpp index 3f7a622..26409c6 100644 --- a/Storage/applications/Client.cpp +++ b/Storage/applications/Client.cpp @@ -1,45 +1,81 @@ -#include -#include -#include -#include -#include -#include +#include "asio.hpp" -constexpr int PORT = 8080; -constexpr int BUFFER_SIZE = 1024; +#include +#include +#include + +// #include +// #include + +constexpr int PORT = 3333; const char* SERVER_IP = "127.0.0.1"; +std::mutex mutex; + +void* readServerBroadcast(void* clientSocketPtr) { + auto socket = ((asio::ip::tcp::socket*)clientSocketPtr); + + while (true) { + + short messageSize; + asio::read(*socket, asio::buffer(&messageSize, 2)); + + mutex.lock(); + + auto message = new char[messageSize + 1]; + message[messageSize] = 0; + + asio::read(*socket, asio::buffer(message, messageSize)); + + std::cerr << "Broadcast : " << message << std::endl; + + delete[] message; + + mutex.unlock(); + } + + return nullptr; +} + int main() { - // Create socket - int clientSocket = socket(AF_INET, SOCK_STREAM, 0); - if (clientSocket == -1) { - std::cerr << "Failed to create socket" << std::endl; + asio::io_context ioContext; + + // Create a TCP socket + asio::ip::tcp::socket socket(ioContext); + + // Connect to a server + asio::ip::tcp::endpoint endpoint(asio::ip::address::from_string(SERVER_IP), PORT); + socket.connect(endpoint); + + // Create a new thread to handle the client + pthread_t threadId; + if (pthread_create(&threadId, nullptr, readServerBroadcast, (void*) &socket) != 0) { + std::cerr << "Failed to create thread for client" << std::endl; return 1; } - // Connect to server - sockaddr_in serverAddress{}; - serverAddress.sin_family = AF_INET; - serverAddress.sin_port = htons(PORT); - if (inet_pton(AF_INET, SERVER_IP, &serverAddress.sin_addr) <= 0) { - std::cerr << "Invalid address or address not supported" << std::endl; - return 1; - } + // Detach the thread so it can run independently + pthread_detach(threadId); - if (connect(clientSocket, reinterpret_cast(&serverAddress), sizeof(serverAddress)) == -1) { - std::cerr << "Failed to connect to server" << std::endl; - return 1; - } + while (true) { + std::string message; + std::cout << " >> "; + std::cin >> message; - // Send message to server - const char* message = "Hello, server!"; - if (write(clientSocket, message, strlen(message)) == -1) { - std::cerr << "Failed to write to socket" << std::endl; - return 1; + mutex.lock(); + + // Send a message to the server + auto messageSize = (short) message.size(); + asio::write(socket, asio::buffer(&messageSize, 2)); + + // Send a message to the server + asio::write(socket, asio::buffer(message + "\n")); + + mutex.unlock(); } // Close socket - close(clientSocket); + // close(clientSocket); return 0; } diff --git a/Storage/applications/Server.cpp b/Storage/applications/Server.cpp index 7f508be..99d9fd3 100644 --- a/Storage/applications/Server.cpp +++ b/Storage/applications/Server.cpp @@ -1,54 +1,56 @@ +#include + +#include #include -#include -#include -#include +#include +#include #include - class Server { + + struct SharedData { + std::list clients; + std::mutex mutex; + }; + + SharedData mSharedData; + + //int serverSocket; + int port; + public: - Server(int port) : serverSocket(-1), port(port) {} + Server(int port) : port(port) {} + ~Server() { assert(0); } bool start() { // Create socket - serverSocket = socket(AF_INET, SOCK_STREAM, 0); - if (serverSocket == -1) { - std::cerr << "Failed to create socket" << std::endl; - return false; - } + asio::io_context io_context; + asio::ip::tcp::acceptor serverSocket(io_context); // Bind socket to port - sockaddr_in serverAddress{}; - serverAddress.sin_family = AF_INET; - serverAddress.sin_port = htons(port); - serverAddress.sin_addr.s_addr = INADDR_ANY; - if (bind(serverSocket, reinterpret_cast(&serverAddress), sizeof(serverAddress)) == -1) { - std::cerr << "Failed to bind socket" << std::endl; - return false; - } + serverSocket.open(asio::ip::tcp::v4()); + // serverSocket.set_option(asio::ip::tcp::acceptor::reuse_address(true)); + serverSocket.bind(asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port)); // Listen for connections - if (listen(serverSocket, 5) == -1) { - std::cerr << "Failed to listen on socket" << std::endl; - return false; - } + serverSocket.listen(); std::cout << "Server listening on port " << port << std::endl; // Start accepting clients while (true) { // Accept client connection - sockaddr_in clientAddress{}; - socklen_t clientAddressSize = sizeof(clientAddress); - int clientSocket = accept(serverSocket, reinterpret_cast(&clientAddress), &clientAddressSize); - if (clientSocket == -1) { - std::cerr << "Failed to accept client connection" << std::endl; - return false; - } + auto clientSocket = new asio::ip::tcp::socket(io_context); + serverSocket.accept(*clientSocket); + + // add client + mSharedData.mutex.lock(); + mSharedData.clients.push_back(clientSocket); + mSharedData.mutex.unlock(); // Create a new thread to handle the client pthread_t threadId; - if (pthread_create(&threadId, nullptr, handleClient, &clientSocket) != 0) { + if (pthread_create(&threadId, nullptr, handleClient, this) != 0) { std::cerr << "Failed to create thread for client" << std::endl; return false; } @@ -58,40 +60,81 @@ public: } // Close server socket - close(serverSocket); + serverSocket.close(); return true; } -private: - int serverSocket; - int port; - static void* handleClient(void* clientSocketPtr) { - int clientSocket = *(reinterpret_cast(clientSocketPtr)); - constexpr int BUFFER_SIZE = 1024; + static void* handleClient(void* in) { + auto self = (Server*) in; - // Receive and print client message - char buffer[BUFFER_SIZE]; - ssize_t bytesRead = read(clientSocket, buffer, BUFFER_SIZE - 1); - if (bytesRead == -1) { - std::cerr << "Failed to read from socket" << std::endl; - return nullptr; + // read shared data - current client id + self->mSharedData.mutex.lock(); + auto clientSocket = self->mSharedData.clients.back(); + self->mSharedData.mutex.unlock(); + + MESSAGE: + // wait for a message request1 + short messageSize; + { + auto bytesRead = asio::read(*clientSocket, asio::buffer(&messageSize, 2)); + if (bytesRead == -1) { + std::cerr << "Failed to read from socket" << std::endl; + (*clientSocket).close(); + return nullptr; + } } - buffer[bytesRead] = '\0'; - std::cout << "Received message from client: " << buffer << std::endl; + self->mSharedData.mutex.lock(); + + // Receive client message + auto message = new char[messageSize + 1]; + message[messageSize] = '\0'; + { + auto bytesRead = asio::read(*clientSocket, asio::buffer(message, messageSize)); + if (bytesRead == -1) { + std::cerr << "Failed to read from socket" << std::endl; + memcpy(message, "Cant wanna say something but i cant read", 100); + } + } + + // Broadcast to all clients + for (auto client: self->mSharedData.clients) { + auto bytesWritten = asio::write(*client, asio::buffer(&messageSize, 2)); + if (bytesWritten == -1) { + std::cerr << "Failed to write to socket" << std::endl; + } + + bytesWritten = write(*client, asio::buffer(message, messageSize)); + if (bytesWritten == -1) { + std::cerr << "Failed to write to socket" << std::endl; + } + } + + auto exit = memcmp(message, "exit", strlen("exit")) == 0; // Close socket - close(clientSocket); + delete[] message; - return nullptr; + + if (exit) { + (*clientSocket).close(); + self->mSharedData.clients.remove(clientSocket); + } + + self->mSharedData.mutex.unlock(); + + if (exit) { + return nullptr; + } + goto MESSAGE; } }; int main() { - Server server(8080); + Server server(3333); server.start(); return 0; -} +} \ No newline at end of file