From 5aec6d6f62373af7d8312e07267be366dbdc96e8 Mon Sep 17 00:00:00 2001 From: Lasan Mahaliyana Date: Sat, 20 Jun 2026 10:52:08 +0530 Subject: [PATCH] added IPC server --- sharing/linux/daemon/BUILD | 20 ++ sharing/linux/daemon/ipc_server.cc | 160 ++++++++++++++ sharing/linux/daemon/ipc_server.h | 44 ++++ sharing/linux/daemon/ipc_server_test.cc | 279 ++++++++++++++++++++++++ sharing/linux/daemon/main.cc | 52 +++++ sharing/linux/daemon/main.h | 45 ++++ 6 files changed, 600 insertions(+) create mode 100644 sharing/linux/daemon/BUILD create mode 100644 sharing/linux/daemon/ipc_server.cc create mode 100644 sharing/linux/daemon/ipc_server.h create mode 100644 sharing/linux/daemon/ipc_server_test.cc create mode 100644 sharing/linux/daemon/main.cc create mode 100644 sharing/linux/daemon/main.h diff --git a/sharing/linux/daemon/BUILD b/sharing/linux/daemon/BUILD new file mode 100644 index 00000000..7fe8b40c --- /dev/null +++ b/sharing/linux/daemon/BUILD @@ -0,0 +1,20 @@ +load("@rules_cc//cc:cc_library.bzl", "cc_library") +load("@rules_cc//cc:cc_test.bzl", "cc_test") + +cc_test( + name = "ipc_server_test", + srcs = ["ipc_server_test.cc"], + deps = [ + ":ipc_server", + "@com_google_googletest//:gtest_main", + "@com_google_absl//absl/synchronization", + ], +) + +cc_library( + name = "ipc_server", + hdrs = ["ipc_server.h"], + deps = [ + "@com_google_absl//absl/synchronization", + ], +) diff --git a/sharing/linux/daemon/ipc_server.cc b/sharing/linux/daemon/ipc_server.cc new file mode 100644 index 00000000..260ecd97 --- /dev/null +++ b/sharing/linux/daemon/ipc_server.cc @@ -0,0 +1,160 @@ +#include "ipc_server.h" +#include +#include +#include +#include + +std::string IPCServer::Read() { + std::string cmd; + absl::MutexLock lock(lock_); + + size_t pos = read_buf.find('\n'); + + if (pos != std::string::npos) { + cmd = read_buf.substr(0, pos); + read_buf.erase(0, pos + 1); + } + + return cmd; +} + +void IPCServer::Stop() { + running_.store(false); + + if (client_fd_ >= 0) { + shutdown(client_fd_, SHUT_RDWR); + close(client_fd_); + client_fd_ = -1; + } + + if (sock_fd_ >= 0) { + shutdown(sock_fd_, SHUT_RDWR); + close(sock_fd_); + sock_fd_ = -1; + } + + unlink(SOCK_PATH.data()); +} + +void IPCServer::InitialiseSock() { + sock_fd_ = socket(AF_UNIX, SOCK_STREAM, 0); + if (sock_fd_ == -1) { + return; + } + + unlink(SOCK_PATH.data()); + + memset(&addr, 0, sizeof(addr)); + addr.sun_family = AF_UNIX; + strncpy(addr.sun_path, SOCK_PATH.data(), sizeof(addr.sun_path) - 1); + + if (bind(sock_fd_, reinterpret_cast(&addr), sizeof(addr)) == -1) { + close(sock_fd_); + sock_fd_ = -1; + return; + } + + if (listen(sock_fd_, 1) == -1) { + close(sock_fd_); + sock_fd_ = -1; + return; + } +} +void IPCServer::Recv() { + while (running_.load()) { + pollfd pfd{}; + pfd.fd = client_fd_; + pfd.events = POLLIN; + + int poll_result = poll(&pfd, 1, 100); + + if (!running_.load()) { + return; + } + + if (poll_result < 0) { + if (errno == EINTR) { + continue; + } + return; + } + + if (poll_result == 0) { + continue; + } + + if (pfd.revents & (POLLNVAL | POLLERR)) { + return; + } + + if (pfd.revents & (POLLIN | POLLHUP)) { + char buf[1024]{}; + + ssize_t n = recv(client_fd_, buf, sizeof(buf), 0); + + if (n > 0) { + absl::MutexLock lock(lock_); + read_buf.append(buf, static_cast(n)); + continue; + } + + if (n == 0) { + return; + } + + if (errno == EINTR) { + continue; + } + + if (errno == EAGAIN || errno == EWOULDBLOCK) { + continue; + } + + return; + } + } +} + +void IPCServer::StartEventLoop() { + running_.store(true); + + InitialiseSock(); + + if (sock_fd_ < 0) { + running_.store(false); + return; + } + + while (running_.load()) { + sockaddr_un client_addr{}; + socklen_t addr_len = sizeof(client_addr); + + int accepted_fd = + accept(sock_fd_, reinterpret_cast(&client_addr), &addr_len); + + if (!running_.load()) { + if (accepted_fd >= 0) { + close(accepted_fd); + } + break; + } + + if (accepted_fd < 0) { + if (errno == EINTR) { + continue; + } + break; + } + + client_fd_ = accepted_fd; + + Recv(); + + if (client_fd_ >= 0) { + close(client_fd_); + client_fd_ = -1; + } + } + + Stop(); +} diff --git a/sharing/linux/daemon/ipc_server.h b/sharing/linux/daemon/ipc_server.h new file mode 100644 index 00000000..2da6cfd9 --- /dev/null +++ b/sharing/linux/daemon/ipc_server.h @@ -0,0 +1,44 @@ +#include +#include +#include + +#include +#include +#include +#include + +#include "absl/synchronization/mutex.h" +#include "gtest/gtest_prod.h" + +constexpr std::string_view SOCK_PATH = "/tmp/nearby_sharing_sock"; + +class IPCServer { + public: + IPCServer() = default; + + ~IPCServer() { Stop(); } + void Stop(); + void InitialiseSock(); + void Recv(); + void StartEventLoop(); + std::string Read(); + + private: + FRIEND_TEST(IPCServerEventLoopTest, ClientCanConnectSendAndServerCanRead); + FRIEND_TEST(IPCServerEventLoopTest, ServerCanReadMultipleCommands); + FRIEND_TEST(IPCServerEventLoopTest, PartialCommandIsBufferedUntilDelimiter); + FRIEND_TEST(IPCServerEventLoopTest, LargeCommandAcrossMultipleRecvCalls); + FRIEND_TEST(IPCServerEventLoopTest, MultipleCommandsSplitAcrossSends); + FRIEND_TEST(IPCServerEventLoopTest, ClientDisconnectDoesNotCrashServer); + FRIEND_TEST(IPCServerEventLoopTest, ClientCanReconnectAfterDisconnect); + FRIEND_TEST(IPCServerEventLoopTest, StopEndsEventLoopThread); + FRIEND_TEST(IPCServerEventLoopTest, StopWhileClientConnected); + FRIEND_TEST(IPCServerEventLoopTest, StaleSocketPathIsCleanedUp); + + int sock_fd_ = -1; + int client_fd_ = -1; + sockaddr_un addr{}; + std::string read_buf; + absl::Mutex lock_; + std::atomic running_{false}; +}; diff --git a/sharing/linux/daemon/ipc_server_test.cc b/sharing/linux/daemon/ipc_server_test.cc new file mode 100644 index 00000000..cc9009e8 --- /dev/null +++ b/sharing/linux/daemon/ipc_server_test.cc @@ -0,0 +1,279 @@ +#include + +#include +#include +#include + +#include +#include +#include +#include +#include + +#include "ipc_server.h" + +namespace { + + +int ConnectClientWithRetry() { + int client_fd = socket(AF_UNIX, SOCK_STREAM, 0); + EXPECT_GE(client_fd, 0); + + sockaddr_un addr{}; + addr.sun_family = AF_UNIX; + strncpy(addr.sun_path, SOCK_PATH.data(), sizeof(addr.sun_path) - 1); + + for (int i = 0; i < 100; ++i) { + int result = connect( + client_fd, + reinterpret_cast(&addr), + sizeof(addr)); + + if (result == 0) { + return client_fd; + } + + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + + ADD_FAILURE() << "connect failed: " << strerror(errno); + close(client_fd); + return -1; +} + +void SendAll(int fd, std::string_view data) { + size_t total_sent = 0; + + while (total_sent < data.size()) { + ssize_t sent = send( + fd, + data.data() + total_sent, + data.size() - total_sent, + 0); + + ASSERT_GT(sent, 0) << "send failed: " << strerror(errno); + + total_sent += static_cast(sent); + } +} + +std::string WaitRead(IPCServer& server) { + for (int i = 0; i < 100; ++i) { + std::string cmd = server.Read(); + + if (!cmd.empty()) { + return cmd; + } + + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + + return ""; +} + + +void StopAndJoin(IPCServer& server, std::thread& server_thread) { + server.Stop(); + + if (server_thread.joinable()) { + server_thread.join(); + } + + unlink(SOCK_PATH.data()); +} + +TEST(IPCServerEventLoopTest, ServerCanReadMultipleCommands) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + SendAll(client_fd, "A\nB\nC\n"); + + EXPECT_EQ(WaitRead(server), "A"); + EXPECT_EQ(WaitRead(server), "B"); + EXPECT_EQ(WaitRead(server), "C"); + + close(client_fd); + StopAndJoin(server, server_thread); +} + +TEST(IPCServerEventLoopTest, PartialCommandIsBufferedUntilDelimiter) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + SendAll(client_fd, "PIN"); + + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + + EXPECT_EQ(server.Read(), ""); + + SendAll(client_fd, "G\n"); + + EXPECT_EQ(WaitRead(server), "PING"); + + close(client_fd); + StopAndJoin(server, server_thread); +} + +TEST(IPCServerEventLoopTest, LargeCommandAcrossMultipleRecvCalls) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + std::string large(8192, 'A'); + SendAll(client_fd, large + "\n"); + + EXPECT_EQ(WaitRead(server), large); + + close(client_fd); + StopAndJoin(server, server_thread); +} + +TEST(IPCServerEventLoopTest, MultipleCommandsSplitAcrossSends) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + SendAll(client_fd, "A\nB"); + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + SendAll(client_fd, "\nC\n"); + + EXPECT_EQ(WaitRead(server), "A"); + EXPECT_EQ(WaitRead(server), "B"); + EXPECT_EQ(WaitRead(server), "C"); + + close(client_fd); + StopAndJoin(server, server_thread); +} + +TEST(IPCServerEventLoopTest, ClientDisconnectDoesNotCrashServer) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + close(client_fd); + + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + + // Main expectation: server did not crash/hang. + StopAndJoin(server, server_thread); + + SUCCEED(); +} + +TEST(IPCServerEventLoopTest, ClientCanReconnectAfterDisconnect) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client1_fd = ConnectClientWithRetry(); + ASSERT_GE(client1_fd, 0); + + SendAll(client1_fd, "FIRST\n"); + EXPECT_EQ(WaitRead(server), "FIRST"); + + close(client1_fd); + + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + + int client2_fd = ConnectClientWithRetry(); + ASSERT_GE(client2_fd, 0); + + SendAll(client2_fd, "SECOND\n"); + EXPECT_EQ(WaitRead(server), "SECOND"); + + close(client2_fd); + StopAndJoin(server, server_thread); +} + +TEST(IPCServerEventLoopTest, StopEndsEventLoopThread) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + + StopAndJoin(server, server_thread); + + SUCCEED(); +} + +TEST(IPCServerEventLoopTest, StopWhileClientConnected) { + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + SendAll(client_fd, "PING\n"); + EXPECT_EQ(WaitRead(server), "PING"); + + StopAndJoin(server, server_thread); + + close(client_fd); + + SUCCEED(); +} + +TEST(IPCServerEventLoopTest, StaleSocketPathIsCleanedUp) { + unlink(SOCK_PATH.data()); + + { + std::ofstream stale_file{SOCK_PATH.data()}; + stale_file << "stale"; + } + + struct stat before{}; + ASSERT_EQ(stat(SOCK_PATH.data(), &before), 0); + ASSERT_TRUE(S_ISREG(before.st_mode)); + + IPCServer server; + + std::thread server_thread([&server]() { + server.StartEventLoop(); + }); + + int client_fd = ConnectClientWithRetry(); + ASSERT_GE(client_fd, 0); + + struct stat after{}; + ASSERT_EQ(stat(SOCK_PATH.data(), &after), 0); + EXPECT_TRUE(S_ISSOCK(after.st_mode)); + + close(client_fd); + StopAndJoin(server, server_thread); +} +} // namespace diff --git a/sharing/linux/daemon/main.cc b/sharing/linux/daemon/main.cc new file mode 100644 index 00000000..2a4ee706 --- /dev/null +++ b/sharing/linux/daemon/main.cc @@ -0,0 +1,52 @@ +#include +#include +#include +#include +#include +#include +#include + +constexpr std::string_view sock_path = "/tmp/nearby_sharing.sock"; + + +int main() { + // should always be waiting on the socket for new connections + int retries = 3; + while (retries >= 0) { + int server_fd; + struct sockaddr_un addr; + + server_fd = socket(AF_UNIX, SOCK_STREAM, 0); + if (server_fd == -1) { + retries--; + continue; + } + + // removes old socket + unlink(sock_path.data()); + + memset(&addr, 0, sizeof(addr)); + addr.sun_family = AF_UNIX; + strncpy(addr.sun_path, sock_path.data(), sizeof(addr.sun_path) - 1); + + if (bind(server_fd, (struct sockaddr*)&addr, sizeof(addr)) == -1) { + close(server_fd); + retries--; + continue; + } + + if (listen(server_fd, 10) == -1) { + close(server_fd); + retries--; + continue; + } + socklen_t addr_len = sizeof(addr); + + // blocks till connection is present + if (accept(server_fd, (struct sockaddr*)&addr, &addr_len)) { + close(server_fd); + retries--; + continue; + } + } +} diff --git a/sharing/linux/daemon/main.h b/sharing/linux/daemon/main.h new file mode 100644 index 00000000..aa7a8ec1 --- /dev/null +++ b/sharing/linux/daemon/main.h @@ -0,0 +1,45 @@ +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include "absl/time/time.h" +#include "connections/implementation/flags/nearby_connections_feature_flags.h" +#include "internal/base/file_path.h" +#include "internal/flags/nearby_flags.h" +#include "sharing/advertisement.h" +#include "sharing/attachment_container.h" +#include "sharing/common/nearby_share_enums.h" +#include "sharing/file_attachment.h" +#include "sharing/linux/platform/linux_sharing_platform.h" +#include "sharing/linux/nearby_noop_analytics_recorder.h" +#include "sharing/flags/generated/nearby_sharing_feature_flags.h" +#include "sharing/nearby_sharing_service.h" +#include "sharing/nearby_sharing_service_factory.h" +#include "sharing/nearby_sharing_settings.h" +#include "sharing/share_target.h" +#include "sharing/share_target_discovered_callback.h" +#include "sharing/transfer_metadata.h" +#include "sharing/transfer_update_callback.h" +#include "sharing/proto/enums.pb.h" + +namespace nearby { +namespace sharing { +namespace linux { + +} +} // namespace sharing +} // namespace nearby +