#include #include "c_TcpServer.h" #if defined(__linux__) #include #include #include #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__ */