From 4f8b46bd57f442fa0ac62e0b05ef45cd261f05db Mon Sep 17 00:00:00 2001 From: "Evgeny Grin (Karlson2k)" Date: Sat, 29 Mar 2025 20:49:14 +0300 Subject: [PATCH] Reworked listen socket/new connections handling --- src/mhd2/daemon_add_conn.c | 1 + src/mhd2/daemon_get_info.c | 19 ++++++--- src/mhd2/daemon_start.c | 9 ++-- src/mhd2/events_process.c | 84 +++++++++++++++++++++++++------------- src/mhd2/mhd_daemon.h | 29 ++++++++++++- 5 files changed, 103 insertions(+), 39 deletions(-) diff --git a/src/mhd2/daemon_add_conn.c b/src/mhd2/daemon_add_conn.c index e2a9e286..2bfc0647 100644 --- a/src/mhd2/daemon_add_conn.c +++ b/src/mhd2/daemon_add_conn.c @@ -793,6 +793,7 @@ mhd_daemon_accept_connection (struct MHD_Daemon *restrict daemon) fd = daemon->net.listen.fd; mhd_assert (MHD_INVALID_SOCKET != fd); + mhd_assert (! daemon->net.listen.is_broken); addrlen = (socklen_t) sizeof (addrstorage); memset (&addrstorage, diff --git a/src/mhd2/daemon_get_info.c b/src/mhd2/daemon_get_info.c index e8585377..412ac7c9 100644 --- a/src/mhd2/daemon_get_info.c +++ b/src/mhd2/daemon_get_info.c @@ -60,7 +60,8 @@ MHD_daemon_get_info_fixed_sz ( switch (info_type) { case MHD_DAEMON_INFO_FIXED_BIND_PORT: - if (MHD_INVALID_SOCKET == daemon->net.listen.fd) + if ((MHD_INVALID_SOCKET == daemon->net.listen.fd) + && ! daemon->net.listen.is_broken) return MHD_SC_INFO_GET_TYPE_NOT_APPLICABLE; if (mhd_SOCKET_TYPE_UNKNOWN > daemon->net.listen.type) return MHD_SC_INFO_GET_TYPE_NOT_APPLICABLE; @@ -75,11 +76,17 @@ MHD_daemon_get_info_fixed_sz ( output_buf->v_bind_port_uint16 = daemon->net.listen.port; return MHD_SC_OK; case MHD_DAEMON_INFO_FIXED_LISTEN_SOCKET: - if (MHD_INVALID_SOCKET == daemon->net.listen.fd) - return MHD_SC_INFO_GET_TYPE_NOT_APPLICABLE; - if (sizeof(output_buf->v_listen_socket) > output_buf_size) - return MHD_SC_INFO_GET_BUFF_TOO_SMALL; - output_buf->v_listen_socket = daemon->net.listen.fd; + if (1) + { + MHD_Socket listen_fd = daemon->net.listen.fd; + if (MHD_INVALID_SOCKET == listen_fd) + return daemon->net.listen.is_broken ? + MHD_SC_INFO_GET_TYPE_UNOBTAINABLE : + MHD_SC_INFO_GET_TYPE_NOT_APPLICABLE; + if (sizeof(output_buf->v_listen_socket) > output_buf_size) + return MHD_SC_INFO_GET_BUFF_TOO_SMALL; + output_buf->v_listen_socket = listen_fd; + } return MHD_SC_OK; case MHD_DAEMON_INFO_FIXED_AGGREAGATE_FD: #ifdef MHD_SUPPORT_EPOLL diff --git a/src/mhd2/daemon_start.c b/src/mhd2/daemon_start.c index d89515a1..0ec3d847 100644 --- a/src/mhd2/daemon_start.c +++ b/src/mhd2/daemon_start.c @@ -645,6 +645,7 @@ create_bind_listen_stream_socket (struct MHD_Daemon *restrict d, { /* No listen socket */ d->net.listen.fd = MHD_INVALID_SOCKET; + d->net.listen.is_broken = false; d->net.listen.type = mhd_SOCKET_TYPE_UNKNOWN; d->net.listen.non_block = false; d->net.listen.port = 0; @@ -988,6 +989,7 @@ create_bind_listen_stream_socket (struct MHD_Daemon *restrict d, /* Set to the daemon only when the listening socket is fully ready */ d->net.listen.fd = sk; + d->net.listen.is_broken = false; switch (sk_type) { case mhd_SKT_UNKNOWN: @@ -1336,7 +1338,6 @@ daemon_choose_and_preinit_events (struct MHD_Daemon *restrict d, mhd_assert ((MHD_WM_EXTERNAL_EVENT_LOOP_CB_LEVEL == s->work_mode.mode) || \ (MHD_WM_EXTERNAL_EVENT_LOOP_CB_EDGE == s->work_mode.mode)); mhd_assert (mhd_WM_INT_HAS_EXT_EVENTS (d->wmode_int)); - mhd_assert (MHD_WM_EXTERNAL_SINGLE_FD_WATCH != s->work_mode.mode); d->events.poll_type = mhd_POLL_TYPE_EXT; d->events.data.ext.cb_data.cb = s->work_mode.params.v_external_event_loop_cb.reg_cb; @@ -1446,7 +1447,7 @@ daemon_init_net (struct MHD_Daemon *restrict d, { if ((MHD_INVALID_SOCKET != d->net.listen.fd) && ! d->net.listen.non_block - && ((mhd_WM_INT_EXTERNAL_EVENTS_EDGE == d->wmode_int) || + && (mhd_D_IS_USING_EDGE_TRIG (d) || (mhd_WM_INT_INTERNAL_EVENTS_THREAD_POOL == d->wmode_int))) { mhd_LOG_MSG (d, MHD_SC_LISTEN_SOCKET_NONBLOCKING_FAILURE, \ @@ -2058,6 +2059,8 @@ init_daemon_fds_monitoring (struct MHD_Daemon *restrict d) mhd_assert (mhd_ITC_IS_VALID (d->threading.itc)); #endif + d->events.accept_pending = false; + switch (d->events.poll_type) { case mhd_POLL_TYPE_EXT: @@ -2168,7 +2171,7 @@ init_daemon_fds_monitoring (struct MHD_Daemon *restrict d) { struct epoll_event reg_event; #ifdef MHD_SUPPORT_THREADS - reg_event.events = EPOLLIN; + reg_event.events = EPOLLIN | EPOLLET; reg_event.data.u64 = (uint64_t) mhd_SOCKET_REL_MARKER_ITC; /* uint64_t is used in the epoll header */ if (0 != epoll_ctl (d->events.data.epoll.e_fd, EPOLL_CTL_ADD, mhd_itc_r_fd (d->threading.itc), ®_event)) diff --git a/src/mhd2/events_process.c b/src/mhd2/events_process.c index 727c45d9..3f22f7d2 100644 --- a/src/mhd2/events_process.c +++ b/src/mhd2/events_process.c @@ -71,10 +71,7 @@ mhd_daemon_get_wait_max (struct MHD_Daemon *restrict d) mhd_assert (! mhd_D_HAS_WORKERS (d)); - if (d->events.act_req.accept && d->conns.block_new) - d->events.act_req.accept = false; - - if (d->events.act_req.accept) + if (d->events.accept_pending && ! d->conns.block_new) return 0; if (zero_wait) return 0; @@ -146,6 +143,7 @@ daemon_accept_new_conns (struct MHD_Daemon *restrict d) { unsigned int num_to_accept; mhd_assert (MHD_INVALID_SOCKET != d->net.listen.fd); + mhd_assert (! d->net.listen.is_broken); mhd_assert (! d->conns.block_new); mhd_assert (d->conns.count < d->conns.cfg.count_limit); mhd_assert (! mhd_D_HAS_WORKERS (d)); @@ -256,9 +254,10 @@ daemon_accept_new_conns (struct MHD_Daemon *restrict d) if (mhd_DAEMON_ACCEPT_NO_MORE_PENDING == res) return true; if (mhd_DAEMON_ACCEPT_FAILED == res) - break; + return false; /* This is probably "no system resources" error. + To do try to accept more connections now. */ } - return false; + return false; /* More connections may need to be accepted */ } @@ -470,8 +469,11 @@ select_update_fdsets (struct MHD_Daemon *restrict d, &ret, d); #endif - if (MHD_INVALID_SOCKET != d->net.listen.fd) + if ((MHD_INVALID_SOCKET != d->net.listen.fd) + && ! d->conns.block_new) { + mhd_assert (! d->net.listen.is_broken); + fd_set_wrap (d->net.listen.fd, rfds, &ret, @@ -544,7 +546,7 @@ select_update_statuses_from_fdsets (struct MHD_Daemon *d, --num_events; /* Clear ITC here, before other data processing. * Any external events will activate ITC again if additional data to - * process is added externally. Cleaning ITC early ensures that new data + * process is added externally. Clearing ITC early ensures that new data * (with additional ITC activation) will not be missed. */ mhd_itc_clear (d->threading.itc); } @@ -557,6 +559,7 @@ select_update_statuses_from_fdsets (struct MHD_Daemon *d, if (MHD_INVALID_SOCKET != d->net.listen.fd) { + mhd_assert (! d->net.listen.is_broken); if (FD_ISSET (d->net.listen.fd, efds)) { --num_events; @@ -567,13 +570,16 @@ select_update_statuses_from_fdsets (struct MHD_Daemon *d, if (! mhd_D_HAS_MASTER (d)) mhd_socket_close (d->net.listen.fd); + d->events.accept_pending = false; + d->net.listen.is_broken = true; /* Stop monitoring socket to avoid spinning with busy-waiting */ d->net.listen.fd = MHD_INVALID_SOCKET; } - else if (FD_ISSET (d->net.listen.fd, rfds)) + else { - --num_events; - d->events.act_req.accept = true; + d->events.accept_pending = FD_ISSET (d->net.listen.fd, rfds); + if (d->events.accept_pending) + --num_events; } } @@ -733,9 +739,11 @@ poll_update_fds (struct MHD_Daemon *restrict d, #endif if (MHD_INVALID_SOCKET != d->net.listen.fd) { + mhd_assert (! d->net.listen.is_broken); mhd_assert (d->events.data.poll.fds[i_s].fd == d->net.listen.fd); mhd_assert (mhd_SOCKET_REL_MARKER_LISTEN == \ d->events.data.poll.rel[i_s].fd_id); + d->events.data.poll.fds[i_s].events = d->conns.block_new ? 0 : POLLIN; ++i_s; } if (listen_only) @@ -809,7 +817,7 @@ poll_update_statuses_from_fds (struct MHD_Daemon *restrict d, --num_events; /* Clear ITC here, before other data processing. * Any external events will activate ITC again if additional data to - * process is added externally. Cleaning ITC early ensures that new data + * process is added externally. Clearing ITC early ensures that new data * (with additional ITC activation) will not be missed. */ mhd_itc_clear (d->threading.itc); } @@ -823,6 +831,7 @@ poll_update_statuses_from_fds (struct MHD_Daemon *restrict d, { const short revents = d->events.data.poll.fds[i_s].revents; + mhd_assert (! d->net.listen.is_broken); mhd_assert (d->events.data.poll.fds[i_s].fd == d->net.listen.fd); mhd_assert (mhd_SOCKET_REL_MARKER_LISTEN == \ d->events.data.poll.rel[i_s].fd_id); @@ -836,13 +845,26 @@ poll_update_statuses_from_fds (struct MHD_Daemon *restrict d, if (! mhd_D_HAS_MASTER (d)) mhd_socket_close (d->net.listen.fd); + d->events.accept_pending = false; + d->net.listen.is_broken = true; /* Stop monitoring socket to avoid spinning with busy-waiting */ d->net.listen.fd = MHD_INVALID_SOCKET; } - else if (0 != (revents & (MHD_POLL_IN | POLLIN))) + else { - --num_events; - d->events.act_req.accept = true; + const bool has_new_conns = (0 != (revents & (MHD_POLL_IN | POLLIN))); + if (has_new_conns) + { + --num_events; + d->events.accept_pending = true; + } + else + { + /* Check whether the listen socket was monitored for incoming + connections */ + if (0 != (d->events.data.poll.fds[i_s].events & POLLIN)) + d->events.accept_pending = false; + } } ++i_s; } @@ -990,7 +1012,7 @@ poll_update_statuses_from_eevents (struct MHD_Daemon *restrict d, { /* Clear ITC here, before other data processing. * Any external events will activate ITC again if additional data to - * process is added externally. Cleaning ITC early ensures that new data + * process is added externally. Clearing ITC early ensures that new data * (with additional ITC activation) will not be missed. */ mhd_itc_clear (d->threading.itc); } @@ -999,23 +1021,32 @@ poll_update_statuses_from_eevents (struct MHD_Daemon *restrict d, #endif /* MHD_SUPPORT_THREADS */ if (((uint64_t) mhd_SOCKET_REL_MARKER_LISTEN) == e->data.u64) /* uint64_t is in the system header */ { + mhd_assert (MHD_INVALID_SOCKET != d->net.listen.fd); if (0 != (e->events & (EPOLLPRI | EPOLLERR | EPOLLHUP))) { mhd_LOG_MSG (d, MHD_SC_LISTEN_STATUS_ERROR, \ "System reported that the listening socket has an error " \ "status. The daemon will not listen any more."); - // TODO: remove listen from epoll monitoring - /* Close the listening socket unless the master daemon should close it */ if (! mhd_D_HAS_MASTER (d)) mhd_socket_close (d->net.listen.fd); + else + { + /* Ignore possible error as the socket could be already removed + from the epoll monitoring by closing the socket */ + (void) epoll_ctl (d->events.data.epoll.e_fd, + EPOLL_CTL_DEL, + d->net.listen.fd, + NULL); + } - /* Stop monitoring socket to avoid spinning with busy-waiting */ + d->events.accept_pending = false; + d->net.listen.is_broken = true; d->net.listen.fd = MHD_INVALID_SOCKET; } - if (0 != (e->events & EPOLLIN)) - d->events.act_req.accept = true; + else + d->events.accept_pending = (0 != (e->events & EPOLLIN)); } else { @@ -1131,13 +1162,10 @@ process_all_events_and_data (struct MHD_Daemon *restrict d) MHD_PANIC ("Daemon data integrity broken"); break; } - if (d->events.act_req.accept) - { - if (daemon_accept_new_conns (d)) - d->events.act_req.accept = false; - else if (! d->net.listen.non_block) - d->events.act_req.accept = false; - } + + if (d->events.accept_pending) + d->events.accept_pending = ! daemon_accept_new_conns (d); + daemon_process_all_active_conns (d); daemon_cleanup_upgraded_conns (d); return ! d->threading.stop_requested; diff --git a/src/mhd2/mhd_daemon.h b/src/mhd2/mhd_daemon.h index ac0dd97f..402d4906 100644 --- a/src/mhd2/mhd_daemon.h +++ b/src/mhd2/mhd_daemon.h @@ -415,9 +415,11 @@ union mhd_DaemonEventMonitoringTypeSpecificData struct mhd_DaemonEventActionRequired { /** - * If 'true' then connection is waiting to be accepted + * If 'true' connection resuming is required. + * TODO: implement resuming + * TODO: add initialisation */ - bool accept; + bool resume; }; @@ -443,6 +445,11 @@ struct mhd_DaemonEventMonitoringData */ struct mhd_DaemonEventActionRequired act_req; + /** + * When set to 'true' the listen socket has new incoming connection(s). + */ + bool accept_pending; + /** * Indicate that daemon already has some data to be processed on the next * cycle @@ -492,6 +499,12 @@ struct mhd_ListenSocket * The listening socket */ MHD_Socket fd; + /** + * Set to 'true' when unrecoverable error is detected on the listen socket. + * If set to 'true', but @a fd is not #MHD_INVALID_SOCKET then "broken" state + * has not yet processed by MHD. + */ + bool is_broken; /** * The type of the listening socket @a fd */ @@ -526,6 +539,12 @@ struct mhd_DaemonNetworkSettings /** * The daemon network/sockets data + * + * This structure holds mostly static information—typically initialised once + * when the daemon starts, and possibly updated once later (e.g., if the + * listening socket fails or is closed). + * + * It does NOT contain any operational states. */ struct mhd_DaemonNetwork { @@ -1060,6 +1079,12 @@ struct MHD_Daemon /** * The daemon network/sockets data + * + * This structure holds mostly static information—typically initialised once + * when the daemon starts, and possibly updated once later (e.g., if the + * listening socket fails or is closed). + * + * It does NOT contain any operational states. */ struct mhd_DaemonNetwork net;