mirror of
https://github.com/kidfromjupiter/nearby.git
synced 2026-09-14 14:46:12 -04:00
added IPC server
This commit is contained in:
@@ -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",
|
||||
],
|
||||
)
|
||||
@@ -0,0 +1,160 @@
|
||||
#include "ipc_server.h"
|
||||
#include <thread>
|
||||
#include <cerrno>
|
||||
#include <cstring>
|
||||
#include <iostream>
|
||||
|
||||
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<sockaddr*>(&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<size_t>(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<sockaddr*>(&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();
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
#include <atomic>
|
||||
#include <string>
|
||||
#include <string_view>
|
||||
|
||||
#include <poll.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/un.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#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<bool> running_{false};
|
||||
};
|
||||
@@ -0,0 +1,279 @@
|
||||
#include <gtest/gtest.h>
|
||||
|
||||
#include <sys/socket.h>
|
||||
#include <sys/un.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include <chrono>
|
||||
#include <cstring>
|
||||
#include <string>
|
||||
#include <fstream>
|
||||
#include <thread>
|
||||
|
||||
#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<sockaddr*>(&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<size_t>(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
|
||||
@@ -0,0 +1,52 @@
|
||||
#include <sys/socket.h>
|
||||
#include <sys/un.h>
|
||||
#include <unistd.h>
|
||||
#include <string>
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <errno.h>
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
#include <signal.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <condition_variable>
|
||||
#include <cstdint>
|
||||
#include <cstdlib>
|
||||
#include <filesystem>
|
||||
#include <functional>
|
||||
#include <iostream>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <optional>
|
||||
#include <string>
|
||||
#include <system_error>
|
||||
#include <utility>
|
||||
#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
|
||||
|
||||
Reference in New Issue
Block a user