From 872068d4b0a376985dbd3acd88da4faa64a57a55 Mon Sep 17 00:00:00 2001 From: wangrong Date: Fri, 11 Sep 2026 15:17:50 +0800 Subject: [PATCH] feat(dispatcher): add event relay dispatcher and receiver modules MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add event_relay_dispatcher and event_relay_receiver for GIO-based DBus event relay, and enable testing in top-level CMakeLists. 新增事件中继调度器与接收器模块,支持通过 DBus 进行事件转发。 Log: 新增事件中继调度器与接收器模块 PMS: BUG-370779 Influence: 新增 dispatcher 事件中继功能,支持 DBus 方式转发事件,并在顶层构建系统中启用 ctest。 --- CMakeLists.txt | 10 + src/dispatcher/CMakeLists.txt | 8 + src/dispatcher/event_relay_dispatcher.c | 131 +++++++ src/dispatcher/event_relay_dispatcher.h | 98 ++++++ src/dispatcher/event_relay_receiver.c | 175 +++++++++ src/dispatcher/event_relay_receiver.h | 91 +++++ src/dispatcher/tests/CMakeLists.txt | 11 + src/dispatcher/tests/test_relay_dispatcher.c | 351 +++++++++++++++++++ 8 files changed, 875 insertions(+) create mode 100644 src/dispatcher/event_relay_dispatcher.c create mode 100644 src/dispatcher/event_relay_dispatcher.h create mode 100644 src/dispatcher/event_relay_receiver.c create mode 100644 src/dispatcher/event_relay_receiver.h create mode 100644 src/dispatcher/tests/test_relay_dispatcher.c diff --git a/CMakeLists.txt b/CMakeLists.txt index 76ded82a..81ff7834 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -8,6 +8,16 @@ set(CMAKE_POSITION_INDEPENDENT_CODE ON) if (CMAKE_INSTALL_PREFIX_INITIALIZED_TO_DEFAULT) set(CMAKE_INSTALL_PREFIX /usr) endif () +# Top-level testing switch: child components (e.g. src/dispatcher) declare +# their own `option(ENABLE_TESTING ... OFF)`, but the cache variable set here +# takes precedence, so a single -DENABLE_TESTING=ON at configure time +# propagates to every component's test subdirectory. +option(ENABLE_TESTING "Enable building tests" OFF) +# Let `ctest` from the top-level build dir traverse into component test +# dirs (each component registers its tests via its own ENABLE_TESTING). +# Must precede add_subdirectory: testing support is only wired into +# subdirectories added after this call. +enable_testing() add_subdirectory(${PROJECT_SOURCE_DIR}/src) add_subdirectory(${PROJECT_SOURCE_DIR}/config) diff --git a/src/dispatcher/CMakeLists.txt b/src/dispatcher/CMakeLists.txt index 92d19e2b..16ab9c4e 100644 --- a/src/dispatcher/CMakeLists.txt +++ b/src/dispatcher/CMakeLists.txt @@ -11,18 +11,26 @@ set(CMAKE_C_STANDARD_REQUIRED ON) find_package(PkgConfig REQUIRED) pkg_check_modules(GLIB REQUIRED glib-2.0) +pkg_check_modules(GIO REQUIRED gio-2.0) +pkg_check_modules(GIO_UNIX REQUIRED gio-unix-2.0) add_library(deepin-anything-dispatcher STATIC event_dispatcher.c event_receiver.c + event_relay_dispatcher.c + event_relay_receiver.c ) target_include_directories(deepin-anything-dispatcher PUBLIC ${GLIB_INCLUDE_DIRS} + ${GIO_INCLUDE_DIRS} + ${GIO_UNIX_INCLUDE_DIRS} ) target_link_libraries(deepin-anything-dispatcher PUBLIC ${GLIB_LIBRARIES} + ${GIO_LIBRARIES} + ${GIO_UNIX_LIBRARIES} ) # Build tests if enabled diff --git a/src/dispatcher/event_relay_dispatcher.c b/src/dispatcher/event_relay_dispatcher.c new file mode 100644 index 00000000..57e928b4 --- /dev/null +++ b/src/dispatcher/event_relay_dispatcher.c @@ -0,0 +1,131 @@ +// SPDX-FileCopyrightText: 2026 UnionTech Software Technology Co., Ltd. +// +// SPDX-License-Identifier: GPL-3.0-or-later + +#define _GNU_SOURCE +#define G_LOG_USE_STRUCTURED +#include "event_relay_dispatcher.h" + +#include +#include +#include +#include +#include +#include + +/* ── Data structure ─────────────────────────────────────────────── */ + +struct EventRelayDispatcher { + GPtrArray *socket_list; /* sender-side fds from socketpair() */ + guint max_event_channel; /* 0 = unlimited */ +}; + +/* ── Helpers ────────────────────────────────────────────────────── */ + +/** + * close_fd: + * @ptr: a #GINT_TO_POINTER-encoded fd + * + * GDestroyNotify-compatible wrapper that closes the fd stored as @ptr. + */ +static void close_fd(gpointer ptr) +{ + int fd = GPOINTER_TO_INT(ptr); + if (fd >= 0) + close(fd); +} + +/* ── Public API ─────────────────────────────────────────────────── */ + +EventRelayDispatcher *event_relay_dispatcher_new(guint max_event_channel) +{ + EventRelayDispatcher *d = g_new0(EventRelayDispatcher, 1); + d->socket_list = g_ptr_array_new_full(8, close_fd); + d->max_event_channel = max_event_channel; + return d; +} + +void event_relay_dispatcher_free(EventRelayDispatcher *dispatcher) +{ + if (dispatcher == NULL) + return; + + /* Closes every fd via the close_fd destroy notify. */ + if (dispatcher->socket_list != NULL) + g_ptr_array_free(dispatcher->socket_list, TRUE); + + g_free(dispatcher); +} + +int event_relay_dispatcher_get_event_channel(EventRelayDispatcher *dispatcher) +{ + g_return_val_if_fail(dispatcher != NULL, -1); + + /* Check channel limit. */ + if (dispatcher->max_event_channel > 0 && + dispatcher->socket_list->len >= dispatcher->max_event_channel) { + g_debug("relay: max_event_channel limit (%u) reached", + dispatcher->max_event_channel); + return -1; + } + + int fds[2]; + if (socketpair(AF_UNIX, SOCK_SEQPACKET | SOCK_NONBLOCK, 0, fds) < 0) { + g_warning("relay: socketpair() failed: %s", strerror(errno)); + return -1; + } + + /* Enlarge send buffer so bursts of events don't cause EAGAIN. */ + int buf_size = 1 << 20; /* 1 MiB */ + if (setsockopt(fds[0], SOL_SOCKET, SO_SNDBUF, &buf_size, sizeof(buf_size)) < 0) + g_debug("relay: setsockopt(SO_SNDBUF) on sender failed: %s", + strerror(errno)); + + /* Store sender end (fds[0]) in the list, return client end (fds[1]). */ + g_ptr_array_add(dispatcher->socket_list, GINT_TO_POINTER(fds[0])); + + g_debug("relay: created event channel (sender fd=%d, client fd=%d, total=%u)", + fds[0], fds[1], dispatcher->socket_list->len); + + return fds[1]; +} + +gboolean event_relay_dispatcher_send(EventRelayDispatcher *dispatcher, + const void *buf, size_t len) +{ + g_return_val_if_fail(dispatcher != NULL, FALSE); + g_return_val_if_fail(buf != NULL, FALSE); + g_return_val_if_fail(len > 0, FALSE); + + /* Iterate in reverse so we can safely remove during iteration. */ + gboolean ok = TRUE; + + for (gint i = dispatcher->socket_list->len - 1; i >= 0; i--) { + int fd = GPOINTER_TO_INT(g_ptr_array_index(dispatcher->socket_list, i)); + + ssize_t ret; + do { + ret = send(fd, buf, len, MSG_NOSIGNAL); + } while (ret < 0 && errno == EINTR); + + if (ret < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK) { + g_debug("relay: slow client fd %d, kicking", fd); + /* Remove from array first so the destroy notify owns the + * close — no double-close, no fd-reuse hazard. */ + g_ptr_array_remove_index_fast(dispatcher->socket_list, i); + } else if (errno == EPIPE || errno == ECONNRESET) { + g_debug("relay: broken connection fd %d", fd); + g_ptr_array_remove_index_fast(dispatcher->socket_list, i); + } else { + g_warning("relay: send() to fd %d failed: %s", + fd, strerror(errno)); + ok = FALSE; + /* Continue to remaining channels — broadcast semantics + * mean one client's error must not starve the others. */ + } + } + } + + return ok; +} diff --git a/src/dispatcher/event_relay_dispatcher.h b/src/dispatcher/event_relay_dispatcher.h new file mode 100644 index 00000000..aec5590e --- /dev/null +++ b/src/dispatcher/event_relay_dispatcher.h @@ -0,0 +1,98 @@ +// SPDX-FileCopyrightText: 2026 UnionTech Software Technology Co., Ltd. +// +// SPDX-License-Identifier: GPL-3.0-or-later + +#ifndef EVENT_RELAY_DISPATCHER_H +#define EVENT_RELAY_DISPATCHER_H + +#include +#include + +G_BEGIN_DECLS + +/** + * EventRelayDispatcher: + * + * An opaque relay event dispatcher. Unlike #EventDispatcher (which filters + * events by per-UID subscription prefixes), the relay dispatcher broadcasts + * every event to every connected channel without filtering. The downstream + * relay agent is responsible for its own filtering. + * + * Communication is via pre-established #socketpair endpoints using + * %SOCK_SEQPACKET (preserves message boundaries). Channel establishment is + * driven by a D-Bus GetEventChannel method call: the + * server-side handler calls event_relay_dispatcher_get_event_channel() to + * create a socketpair, stores one end, and returns the other end to the + * client via D-Bus fd passing. + */ +typedef struct EventRelayDispatcher EventRelayDispatcher; + +/** + * event_relay_dispatcher_new: (constructor) + * @max_event_channel: Maximum number of concurrent event channels (0 = unlimited) + * + * Creates a new relay event dispatcher. The dispatcher maintains a list of + * socketpair endpoints (one per channel). When a client requests a channel + * via event_relay_dispatcher_get_event_channel(), a new socketpair is created + * and one end is stored in the list while the other is returned to the caller. + * + * Returns: (transfer full) (nullable): a new #EventRelayDispatcher, or %NULL + * on failure + */ +EventRelayDispatcher *event_relay_dispatcher_new(guint max_event_channel); + +/** + * event_relay_dispatcher_free: (skip) + * @dispatcher: (nullable) + * + * Frees the relay dispatcher and all associated resources. Closes every + * socketpair endpoint in the list. Safe to call on %NULL. + * + * Not thread-safe relative to get_event_channel/send — the caller must + * ensure no concurrent calls. + */ +void event_relay_dispatcher_free(EventRelayDispatcher *dispatcher); + +/** + * event_relay_dispatcher_get_event_channel: + * @dispatcher: an #EventRelayDispatcher + * + * Creates a new event channel by calling socketpair(%AF_UNIX, + * %SOCK_SEQPACKET | %SOCK_NONBLOCK, 0). One end of the pair is stored + * in the dispatcher's internal socket list; the other end is returned + * as a file descriptor. + * + * If @max_event_channel is non-zero and the number of active channels + * has already reached the limit, this function returns -1. + * + * The D-Bus service framework calls this function to respond to a + * GetEventChannel method call. The returned fd is + * passed back to the client via D-Bus fd passing. The service framework + * also returns an @event_protocol_id that identifies the data structure + * of the buffer used in send() — the current protocol ID is 0, meaning + * #dispatch_event_t. + * + * Returns: a file descriptor (the client end of the socketpair), or -1 + * on failure or if the channel limit has been reached + */ +int event_relay_dispatcher_get_event_channel(EventRelayDispatcher *dispatcher); + +/** + * event_relay_dispatcher_send: + * @dispatcher: an #EventRelayDispatcher + * @buf: the data to send + * @len: the length of @buf in bytes + * + * Sends @buf to every active channel in the dispatcher's socket list. + * Disconnected or slow clients (EAGAIN/EPIPE/ECONNRESET on send) are + * automatically cleaned up — their fd is closed and removed from the list. + * + * Returns: %TRUE if the data was processed (even if no clients received it), + * %FALSE on critical error + */ +gboolean event_relay_dispatcher_send(EventRelayDispatcher *dispatcher, + const void *buf, size_t len); + +G_END_DECLS + +#endif /* EVENT_RELAY_DISPATCHER_H */ diff --git a/src/dispatcher/event_relay_receiver.c b/src/dispatcher/event_relay_receiver.c new file mode 100644 index 00000000..c0b01d3a --- /dev/null +++ b/src/dispatcher/event_relay_receiver.c @@ -0,0 +1,175 @@ +// SPDX-FileCopyrightText: 2026 UnionTech Software Technology Co., Ltd. +// +// SPDX-License-Identifier: GPL-3.0-or-later + +#define _GNU_SOURCE +#define G_LOG_USE_STRUCTURED +#include "event_relay_receiver.h" + +#include +#include +#include +#include +#include + +#include +#include + +/* 1 MiB receive buffer — same rationale as the existing EventReceiver. */ +#define RELAY_RECEIVER_SOCKET_BUF_SIZE (1 << 20) + +/* ── Data structure ─────────────────────────────────────────────── */ + +struct EventRelayReceiver { + int sock_fd; /* socket fd from GetEventChannel */ + guint32 event_protocol_id; /* protocol ID returned by GetEventChannel */ +}; + +/* ── Public API ─────────────────────────────────────────────────── */ + +EventRelayReceiver *event_relay_receiver_new(const char *bus_name, + const char *object_path, + const char *interface_name) +{ + g_return_val_if_fail(bus_name != NULL, NULL); + g_return_val_if_fail(object_path != NULL, NULL); + g_return_val_if_fail(interface_name != NULL, NULL); + + EventRelayReceiver *r = g_new0(EventRelayReceiver, 1); + r->sock_fd = -1; + r->event_protocol_id = 0; + + GError *error = NULL; + + /* Connect to the system bus (server is a root system service). */ + GDBusConnection *connection = g_bus_get_sync(G_BUS_TYPE_SYSTEM, NULL, &error); + if (connection == NULL) { + g_warning("relay: failed to connect to system bus: %s", + error ? error->message : "unknown error"); + g_clear_error(&error); + goto fail; + } + + /* Call GetEventChannel on the specified D-Bus service. The method + * returns (h fd, u event_protocol_id) with a GUnixFDList for fd passing. */ + GUnixFDList *fd_list = NULL; + GVariant *result = g_dbus_connection_call_with_unix_fd_list_sync( + connection, + bus_name, + object_path, + interface_name, + "GetEventChannel", + NULL, /* parameters — none */ + G_VARIANT_TYPE("(hu)"), + G_DBUS_CALL_FLAGS_NONE, + -1, /* timeout: default */ + NULL, /* fd_list to send — none */ + &fd_list, + NULL, /* cancellable */ + &error); + + g_object_unref(connection); + + if (result == NULL) { + g_warning("relay: GetEventChannel call failed: %s", + error ? error->message : "unknown error"); + g_clear_error(&error); + goto fail; + } + + /* Extract the fd handle and event_protocol_id from the return value. + * The D-Bus signature is (h u): h is the index into the GUnixFDList + * for the passed fd, u is the event_protocol_id. */ + guint32 protocol_id = 0; + guint32 fd_handle = 0; + g_variant_get(result, "(hu)", &fd_handle, &protocol_id); + g_variant_unref(result); + + if (fd_list == NULL || g_unix_fd_list_get_length(fd_list) < 1) { + g_warning("relay: GetEventChannel returned no file descriptor"); + g_clear_object(&fd_list); + goto fail; + } + + if (fd_handle >= (guint32)g_unix_fd_list_get_length(fd_list)) { + g_warning("relay: GetEventChannel fd handle %u out of range (len=%d)", + fd_handle, g_unix_fd_list_get_length(fd_list)); + g_clear_object(&fd_list); + goto fail; + } + + g_clear_error(&error); + gint fd = g_unix_fd_list_get(fd_list, (gint)fd_handle, &error); + g_object_unref(fd_list); + + if (fd < 0) { + g_warning("relay: GetEventChannel returned invalid fd: %s", + error ? error->message : "unknown error"); + g_clear_error(&error); + goto fail; + } + + r->sock_fd = fd; + r->event_protocol_id = protocol_id; + + /* Enlarge the receive buffer so bursts don't cause the server to kick us. */ + int buf_size = RELAY_RECEIVER_SOCKET_BUF_SIZE; + if (setsockopt(r->sock_fd, SOL_SOCKET, SO_RCVBUF, &buf_size, + sizeof(buf_size)) < 0) + g_debug("relay: setsockopt(SO_RCVBUF) failed: %s", strerror(errno)); + + g_debug("relay: receiver created (fd=%d, protocol_id=%u)", + r->sock_fd, r->event_protocol_id); + + return r; + +fail: + event_relay_receiver_free(r); + return NULL; +} + +void event_relay_receiver_free(EventRelayReceiver *receiver) +{ + if (receiver == NULL) + return; + + if (receiver->sock_fd >= 0) + close(receiver->sock_fd); + + g_free(receiver); +} + +gboolean event_relay_receiver_get_fd(EventRelayReceiver *receiver, + int *fd, guint32 *event_protocol_id) +{ + g_return_val_if_fail(receiver != NULL, FALSE); + + if (fd != NULL) + *fd = receiver->sock_fd; + if (event_protocol_id != NULL) + *event_protocol_id = receiver->event_protocol_id; + + return TRUE; +} + +EventReceiveResult event_relay_receiver_receive(EventRelayReceiver *receiver, + void *buf, size_t len) +{ + g_return_val_if_fail(receiver != NULL, EVENT_RECEIVE_ERROR); + g_return_val_if_fail(buf != NULL, EVENT_RECEIVE_ERROR); + + ssize_t ret = recv(receiver->sock_fd, buf, len, MSG_DONTWAIT); + if (ret < 0) { + if (errno == EINTR) + return EVENT_RECEIVE_INTERRUPTED; + if (errno == EAGAIN || errno == EWOULDBLOCK) + return EVENT_RECEIVE_WOULD_BLOCK; + g_debug("relay: recv() failed: %s", strerror(errno)); + return EVENT_RECEIVE_ERROR; + } + + if (ret == 0) + return EVENT_RECEIVE_DISCONNECTED; + + return EVENT_RECEIVE_OK; +} diff --git a/src/dispatcher/event_relay_receiver.h b/src/dispatcher/event_relay_receiver.h new file mode 100644 index 00000000..24902a80 --- /dev/null +++ b/src/dispatcher/event_relay_receiver.h @@ -0,0 +1,91 @@ +// SPDX-FileCopyrightText: 2026 UnionTech Software Technology Co., Ltd. +// +// SPDX-License-Identifier: GPL-3.0-or-later + +#ifndef EVENT_RELAY_RECEIVER_H +#define EVENT_RELAY_RECEIVER_H + +#include +#include + +#include "event_dispatcher.h" + +G_BEGIN_DECLS + +/** + * EventRelayReceiver: + * + * An opaque relay event receiver. It obtains a communication socket from a + * D-Bus service's GetEventChannel method and reads events + * from it. Unlike #EventReceiver (which connects to a Unix Domain Socket + * path), the relay receiver gets its fd via D-Bus fd passing. + */ + +typedef struct EventRelayReceiver EventRelayReceiver; + +/** + * event_relay_receiver_new: (constructor) + * @bus_name: the D-Bus bus name of the service to connect to + * @object_path: the D-Bus object path of the service + * @interface_name: the D-Bus interface name providing GetEventChannel + * + * Creates a new relay event receiver by calling the specified D-Bus + * service's GetEventChannel method. The method returns + * a file descriptor (via D-Bus fd passing) and an @event_protocol_id + * (uint32) that identifies the data structure of the event buffer. + * Both are stored in the receiver for later use. + * + * Returns: (transfer full) (nullable): a new #EventRelayReceiver, or %NULL + * on failure + */ +EventRelayReceiver *event_relay_receiver_new(const char *bus_name, + const char *object_path, + const char *interface_name); + +/** + * event_relay_receiver_free: (skip) + * @receiver: (nullable) + * + * Frees the relay receiver and closes the socket fd. Safe to call on %NULL. + */ +void event_relay_receiver_free(EventRelayReceiver *receiver); + +/** + * event_relay_receiver_get_fd: + * @receiver: an #EventRelayReceiver + * @fd: (out) (optional): location to store the socket fd, or %NULL + * @event_protocol_id: (out) (optional): location to store the protocol ID, or %NULL + * + * Retrieves the socket fd and @event_protocol_id from the receiver. The + * caller can monitor the fd for readability (e.g., via poll/epoll) and + * call event_relay_receiver_receive() when data is available. + * + * Returns: %TRUE on success, %FALSE if @receiver is %NULL + */ +gboolean event_relay_receiver_get_fd(EventRelayReceiver *receiver, + int *fd, guint32 *event_protocol_id); + +/** + * event_relay_receiver_receive: + * @receiver: an #EventRelayReceiver + * @buf: (out caller-allocates): output buffer + * @len: the maximum length of @buf in bytes + * + * Non-blocking variant of receive. Uses %MSG_DONTWAIT so the call returns + * immediately when no data is available. Designed for use in fd-readable + * callbacks that drain the socket in a loop. %SOCK_SEQPACKET preserves + * message boundaries, so a single recv() returns one complete message. + * + * Returns: #EventReceiveResult indicating the outcome: + * %EVENT_RECEIVE_OK on success, + * %EVENT_RECEIVE_WOULD_BLOCK when no data is available, + * %EVENT_RECEIVE_DISCONNECTED when the server closed the connection, + * %EVENT_RECEIVE_INTERRUPTED when recv() was interrupted by a signal, + * %EVENT_RECEIVE_ERROR on other errors + */ +EventReceiveResult event_relay_receiver_receive(EventRelayReceiver *receiver, + void *buf, size_t len); + +G_END_DECLS + +#endif /* EVENT_RELAY_RECEIVER_H */ diff --git a/src/dispatcher/tests/CMakeLists.txt b/src/dispatcher/tests/CMakeLists.txt index e8dd1c16..51973876 100644 --- a/src/dispatcher/tests/CMakeLists.txt +++ b/src/dispatcher/tests/CMakeLists.txt @@ -14,3 +14,14 @@ target_include_directories(test_dispatcher PRIVATE ${GLIB_INCLUDE_DIRS} ) add_test(NAME dispatcher COMMAND test_dispatcher) + +add_executable(test_relay_dispatcher test_relay_dispatcher.c) +target_link_libraries(test_relay_dispatcher PRIVATE + deepin-anything-dispatcher + ${GLIB_LIBRARIES} +) +target_include_directories(test_relay_dispatcher PRIVATE + ${CMAKE_CURRENT_SOURCE_DIR}/.. + ${GLIB_INCLUDE_DIRS} +) +add_test(NAME relay-dispatcher COMMAND test_relay_dispatcher) diff --git a/src/dispatcher/tests/test_relay_dispatcher.c b/src/dispatcher/tests/test_relay_dispatcher.c new file mode 100644 index 00000000..bf2a0bd4 --- /dev/null +++ b/src/dispatcher/tests/test_relay_dispatcher.c @@ -0,0 +1,351 @@ +// SPDX-FileCopyrightText: 2026 UnionTech Software Technology Co., Ltd. +// +// SPDX-License-Identifier: GPL-3.0-or-later + +#define _GNU_SOURCE +#define G_LOG_USE_STRUCTURED +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "event_dispatcher.h" +#include "event_relay_dispatcher.h" +#include "event_relay_receiver.h" + +/* ── Sender tests (no D-Bus needed) ────────────────────────────── */ + +static void test_relay_new_free(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + event_relay_dispatcher_free(d); +} + +static void test_relay_new_free_null(void) +{ + event_relay_dispatcher_free(NULL); +} + +static void test_relay_get_event_channel(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int fd = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd, >=, 0); + + /* fd should be non-blocking */ + int flags = fcntl(fd, F_GETFL, 0); + g_assert_true(flags & O_NONBLOCK); + + close(fd); + event_relay_dispatcher_free(d); +} + +static void test_relay_max_channel_limit(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(2); + g_assert_nonnull(d); + + int fd1 = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd1, >=, 0); + + int fd2 = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd2, >=, 0); + + /* Third channel should be rejected (limit = 2) */ + int fd3 = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd3, ==, -1); + + close(fd1); + close(fd2); + event_relay_dispatcher_free(d); +} + +static void test_relay_send_no_clients(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + dispatch_event_t event = {0}; + event.event_action = 1; + g_strlcpy(event.event_path, "/test", sizeof(event.event_path)); + + gboolean ret = event_relay_dispatcher_send(d, &event, sizeof(event)); + g_assert_true(ret); + + event_relay_dispatcher_free(d); +} + +static void test_relay_send_receive(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int client_fd = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(client_fd, >=, 0); + + /* Small delay to ensure socketpair is fully set up */ + usleep(10000); + + dispatch_event_t send_event = {0}; + send_event.event_action = 42; + g_strlcpy(send_event.event_path, "/home/user/relay_test.txt", + sizeof(send_event.event_path)); + + gboolean ret = event_relay_dispatcher_send(d, &send_event, sizeof(send_event)); + g_assert_true(ret); + + /* Receive on the client end — use poll for non-blocking wait */ + struct pollfd pfd; + pfd.fd = client_fd; + pfd.events = POLLIN; + g_assert_cmpint(poll(&pfd, 1, 5000), >, 0); + g_assert_true(pfd.revents & POLLIN); + + dispatch_event_t recv_event = {0}; + ssize_t n = recv(client_fd, &recv_event, sizeof(recv_event), 0); + g_assert_cmpint(n, ==, sizeof(recv_event)); + g_assert_cmpint(recv_event.event_action, ==, 42); + g_assert_cmpstr(recv_event.event_path, ==, "/home/user/relay_test.txt"); + + close(client_fd); + event_relay_dispatcher_free(d); +} + +static void test_relay_send_multiple_clients(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int fd1 = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd1, >=, 0); + int fd2 = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(fd2, >=, 0); + + usleep(10000); + + dispatch_event_t event = {0}; + event.event_action = 7; + g_strlcpy(event.event_path, "/tmp/multi.txt", sizeof(event.event_path)); + + gboolean ret = event_relay_dispatcher_send(d, &event, sizeof(event)); + g_assert_true(ret); + + /* Both clients should receive the event */ + for (int i = 0; i < 2; i++) { + int fd = (i == 0) ? fd1 : fd2; + struct pollfd pfd; + pfd.fd = fd; + pfd.events = POLLIN; + g_assert_cmpint(poll(&pfd, 1, 5000), >, 0); + + dispatch_event_t recv_event = {0}; + ssize_t n = recv(fd, &recv_event, sizeof(recv_event), 0); + g_assert_cmpint(n, ==, sizeof(recv_event)); + g_assert_cmpint(recv_event.event_action, ==, 7); + g_assert_cmpstr(recv_event.event_path, ==, "/tmp/multi.txt"); + } + + close(fd1); + close(fd2); + event_relay_dispatcher_free(d); +} + +static void test_relay_cleanup_disconnected(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int client_fd = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(client_fd, >=, 0); + + /* Close the client end — simulates disconnection */ + close(client_fd); + + /* Send should clean up the dead channel without crashing */ + dispatch_event_t event = {0}; + event.event_action = 1; + g_strlcpy(event.event_path, "/test", sizeof(event.event_path)); + + gboolean ret = event_relay_dispatcher_send(d, &event, sizeof(event)); + g_assert_true(ret); + + /* After cleanup, send again — should still work (empty list) */ + ret = event_relay_dispatcher_send(d, &event, sizeof(event)); + g_assert_true(ret); + + event_relay_dispatcher_free(d); +} + +static void test_relay_send_partial_msg(void) +{ + /* Test sending a partial dispatch_event_t (only the header portion). + * SOCK_SEQPACKET preserves message boundaries so recv gets exactly + * what was sent. */ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int client_fd = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(client_fd, >=, 0); + + usleep(10000); + + /* Send only the first 8 bytes (event_action + cookie, no path) */ + struct { + gint32 action; + guint32 cookie; + } partial = { .action = 99, .cookie = 42 }; + + gboolean ret = event_relay_dispatcher_send(d, &partial, sizeof(partial)); + g_assert_true(ret); + + struct pollfd pfd; + pfd.fd = client_fd; + pfd.events = POLLIN; + g_assert_cmpint(poll(&pfd, 1, 5000), >, 0); + + char buf[sizeof(partial)]; + ssize_t n = recv(client_fd, buf, sizeof(buf), 0); + g_assert_cmpint(n, ==, sizeof(partial)); + + gint32 *action = (gint32 *)buf; + guint32 *cookie = (guint32 *)(buf + 4); + g_assert_cmpint(*action, ==, 99); + g_assert_cmpuint(*cookie, ==, 42); + + close(client_fd); + event_relay_dispatcher_free(d); +} + +/* ── Receiver tests (no D-Bus — test receive logic directly) ────── */ + +/* The receiver's new() requires a D-Bus service, which is not available + * in unit tests. However, we can test the receive logic by creating a + * socketpair, pretending one end is the receiver's fd, and calling + * event_relay_receiver_receive() — but that requires constructing the + * opaque struct. Since the struct is opaque, we test the receive logic + * indirectly via the sender tests (which already exercise send→recv + * on the socketpair). For the receiver D-Bus path, we add a test that + * verifies event_relay_receiver_new returns NULL when no D-Bus service + * is available. */ + +static void test_relay_receiver_new_no_dbus(void) +{ + /* With a non-existent bus name, the D-Bus call should fail and the + * constructor should return NULL. We fork so that the g_warning + * (which would be fatal in the parent under + * g_log_set_always_fatal) only affects the child. */ + pid_t pid = fork(); + g_assert_cmpint(pid, >=, 0); + + if (pid == 0) { + /* Child: disable fatal warnings so we can reach the return path + * and verify it directly instead of relying on a signal abort. */ + g_log_set_always_fatal(0); + + EventRelayReceiver *r = event_relay_receiver_new( + "com.test.NonExistentService", + "/com/test/NonExistent", + "com.test.NonExistentInterface"); + + /* Constructor must return NULL on D-Bus failure. */ + event_relay_receiver_free(r); + _exit(r == NULL ? 0 : 1); + } + + /* Parent: require a normal exit with status 0. We do NOT accept + * arbitrary signal termination as success — a crash (segfault, + * abort from an unrelated cause) is a real bug and must fail the + * test rather than being masked as a "expected fatal warning". */ + int status; + waitpid(pid, &status, 0); + + g_assert_true(WIFEXITED(status)); + g_assert_cmpint(WEXITSTATUS(status), ==, 0); +} + +/* ── End-to-end fork test ───────────────────────────────────────── */ + +static void test_relay_end_to_end_fork(void) +{ + EventRelayDispatcher *d = event_relay_dispatcher_new(0); + g_assert_nonnull(d); + + int client_fd = event_relay_dispatcher_get_event_channel(d); + g_assert_cmpint(client_fd, >=, 0); + + pid_t pid = fork(); + g_assert_cmpint(pid, >=, 0); + + if (pid == 0) { + /* Child: receive the event */ + dispatch_event_t event = {0}; + + struct pollfd pfd; + pfd.fd = client_fd; + pfd.events = POLLIN; + + int nfds = poll(&pfd, 1, 5000); + if (nfds <= 0) + _exit(1); + + ssize_t n = recv(client_fd, &event, sizeof(event), 0); + if (n != sizeof(event)) + _exit(1); + if (event.event_action != 55) + _exit(2); + if (strcmp(event.event_path, "/e2e/fork/test") != 0) + _exit(3); + + close(client_fd); + _exit(0); + } + + /* Parent: send the event */ + usleep(10000); + + dispatch_event_t event = {0}; + event.event_action = 55; + g_strlcpy(event.event_path, "/e2e/fork/test", sizeof(event.event_path)); + + gboolean ret = event_relay_dispatcher_send(d, &event, sizeof(event)); + g_assert_true(ret); + + int status; + waitpid(pid, &status, 0); + g_assert_true(WIFEXITED(status)); + g_assert_cmpint(WEXITSTATUS(status), ==, 0); + + close(client_fd); + event_relay_dispatcher_free(d); +} + +/* ── Main ──────────────────────────────────────────────────────── */ + +int main(int argc, char *argv[]) +{ + g_test_init(&argc, &argv, NULL); + g_log_set_always_fatal(G_LOG_LEVEL_CRITICAL | G_LOG_LEVEL_WARNING); + + g_test_add_func("/relay/new_free", test_relay_new_free); + g_test_add_func("/relay/new_free_null", test_relay_new_free_null); + g_test_add_func("/relay/get_event_channel", test_relay_get_event_channel); + g_test_add_func("/relay/max_channel_limit", test_relay_max_channel_limit); + g_test_add_func("/relay/send_no_clients", test_relay_send_no_clients); + g_test_add_func("/relay/send_receive", test_relay_send_receive); + g_test_add_func("/relay/send_multiple_clients", test_relay_send_multiple_clients); + g_test_add_func("/relay/cleanup_disconnected", test_relay_cleanup_disconnected); + g_test_add_func("/relay/send_partial_msg", test_relay_send_partial_msg); + g_test_add_func("/relay/receiver_new_no_dbus", test_relay_receiver_new_no_dbus); + g_test_add_func("/relay/end_to_end_fork", test_relay_end_to_end_fork); + + return g_test_run(); +}