Files
cKit/Foundation/c_TcpEpollReactor.c
2026-09-08 12:38:28 +08:00

175 lines
6.6 KiB
C

#include <c_TcpEpollReactor.h>
#include "c_TcpServer.h"
#if defined(__linux__)
#include <sys/epoll.h>
#include <unistd.h>
#include <string.h>
#define MAX_EPOLL_EVENTS 1024
typedef struct {
c_TcpReactor_t base;
int epoll_fd; /* Raw Linux epoll kernel handle instance descriptor */
struct epoll_event events[MAX_EPOLL_EVENTS]; /* Pre-allocated active kernel alerts buffer */
c_Allocator_t allocator; /* Deep copy of the user-provided allocator */
} c_EPollReactorImpl_t;
static c_err_t EPoll_Add(c_TcpReactor_t* self, c_socket_t sock, uint32_t events) {
c_EPollReactorImpl_t* impl = (c_EPollReactorImpl_t*)self;
struct epoll_event ev;
memset(&ev, 0, sizeof(ev));
/* Map incoming events to POLLIN / read data availability triggers */
ev.events = (events ? EPOLLIN : 0) | EPOLLERR | EPOLLHUP;
ev.data.fd = (int)sock;
if (epoll_ctl(impl->epoll_fd, EPOLL_CTL_ADD, (int)sock, &ev) < 0) {
return C_ERR_FAIL;
}
return C_SUCCESS;
}
static c_err_t EPoll_Remove(c_TcpReactor_t* self, c_socket_t sock) {
c_EPollReactorImpl_t* impl = (c_EPollReactorImpl_t*)self;
/* Passing NULL for the epoll_event structure is valid for deletions since Linux 2.6.9 */
if (epoll_ctl(impl->epoll_fd, EPOLL_CTL_DEL, (int)sock, NULL) < 0) {
return C_ERR_NOTFOUND;
}
return C_SUCCESS;
}
static c_err_t EPoll_Poll(c_TcpReactor_t* self, int timeout_ms, c_TcpServer_t* server, void* args) {
c_EPollReactorImpl_t* impl = (c_EPollReactorImpl_t*)self;
if (!server->is_running) return C_SUCCESS;
/* Intercept readiness frames directly out of the epoll kernel tree instance */
int ret = epoll_wait(impl->epoll_fd, impl->events, MAX_EPOLL_EVENTS, timeout_ms);
if (ret < 0) {
if (server->fnOnError) {
server->fnOnError(server, server->listen_sock, C_TCPSERVER_ERR_ON_POLL, args);
}
return C_SUCCESS;
}
if (ret == 0) return C_SUCCESS; /* Timeout pass, exit cleanly */
/*
* Iterative Forward Processing Pass:
* Unlike arrays, epoll_wait packs active ready elements sequentially into indices 0 to ret-1.
* This provides a strict O(1) loop traversal path that is entirely immune to index mutations.
*/
for (int i = 0; i < ret; ++i) {
struct epoll_event* current_ev = &impl->events[i];
c_socket_t active_sock = (c_socket_t)current_ev->data.fd;
/* Channel A: Master server listener handles incoming connection handshakes */
if (active_sock == server->listen_sock) {
if (current_ev->events & EPOLLIN) {
c_SockAddr_t peer_addr;
socklen_t addr_len = sizeof(peer_addr);
c_socket_t client = accept(server->listen_sock, (struct sockaddr*)&peer_addr, &addr_len);
if (c_Socket_IsValid(client)) {
bool keep = true;
if (server->fnOnConnect) {
keep = server->fnOnConnect(server, client, &peer_addr, args);
}
if (keep) {
EPoll_Add(self, client, 1); /* Inject new client into the red-black tree watchlist */
} else {
c_Socket_Close(client);
}
} else if (server->fnOnError) {
server->fnOnError(server, server->listen_sock, C_TCPSERVER_ERR_ON_ACCEPT, args);
}
}
if (current_ev->events & (EPOLLERR | EPOLLHUP)) {
if (server->fnOnError) {
server->fnOnError(server, server->listen_sock, C_TCPSERVER_ERR_ON_POLLERR, args);
}
}
}
/* Channel B: Connected active client descriptors processing data readiness loops */
else {
c_bool_t keep_alive = C_TRUE;
/* Intercept hardware-level operational error signals first */
if (current_ev->events & EPOLLERR) {
if (server->fnOnError) {
server->fnOnError(server, active_sock, C_TCPSERVER_ERR_ON_POLLERR, args);
}
keep_alive = C_FALSE;
}
if (keep_alive && (current_ev->events & EPOLLIN)) {
/* Identify clean socket EOF drops via zero-byte peeks */
char peek_buf;
ssize_t peek_res = recv(active_sock, &peek_buf, 1, MSG_PEEK);
if (peek_res == 0 || peek_res < 0) {
keep_alive = C_FALSE; /* Connection explicitly severed by remote peer node */
} else if (server->fnOnRequest) {
keep_alive = server->fnOnRequest(server, active_sock, args);
}
}
/* Clean up active client descriptor tracking states if loop signaling drops or errors out */
if (!keep_alive || (current_ev->events & EPOLLHUP)) {
/*
* Remove explicitly from epoll tracking first.
* Closing the socket handle automatically removes it from epoll,
* but manual synchronization ensures clean deterministic state tracking.
*/
EPoll_Remove(self, active_sock);
if (server->fnOnDisconnect) {
server->fnOnDisconnect(server, active_sock, args);
}
c_Socket_Close(active_sock);
}
}
}
return C_SUCCESS;
}
static void EPoll_Destroy(c_TcpReactor_t* self) {
c_EPollReactorImpl_t* impl = (c_EPollReactorImpl_t*)self;
c_Allocator_t alloc = impl->allocator;
if (impl->epoll_fd >= 0) {
close(impl->epoll_fd);
}
c_Allocator_Free(&alloc, impl);
}
static const c_TcpReactorVtbl_t g_EPollReactorVtbl = { EPoll_Add, EPoll_Remove, EPoll_Poll, EPoll_Destroy };
/* ------------------------------------------------------------------------------------------------------------------ */
/* */
c_err_t c_TcpEPollReactor_Create(c_TcpReactor_t** out_reactor, c_Allocator_t* allocator) {
if (!out_reactor) return C_ERR_PARAM;
c_Allocator_t alloc = allocator ? *allocator : c_DefaultAllocator;
c_EPollReactorImpl_t* impl = (c_EPollReactorImpl_t*)c_Allocator_Calloc(&alloc, 1, sizeof(c_EPollReactorImpl_t));
if (!impl) return C_ERR_NOMEM;
/* Create the epoll instance handle (Size parameter must be > 0; ignored since Linux 2.6.8) */
impl->epoll_fd = epoll_create(1);
if (impl->epoll_fd < 0) {
c_Allocator_Free(&alloc, impl);
return C_ERR_FAIL;
}
impl->base.vtbl = &g_EPollReactorVtbl;
impl->allocator = alloc;
*out_reactor = &impl->base;
return C_SUCCESS;
}
#endif /* __linux__ */