From b43797bd4c5daa8134910d4eda5f013f0fa4eb2d Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sat, 25 Jul 2026 19:42:52 +0000 Subject: [PATCH 01/15] [dhcpmon]: Add packet event quiescing Allow packet handlers to be suspended and resumed without terminating socket event loops, and bound each callback batch so quiescing completes under sustained traffic. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 21 ++++++++++++++++++++- src/event_mgr.h | 2 ++ src/packet_handler.cpp | 10 ++++++++-- src/sock_mgr.cpp | 35 +++++++++++++++++++++++++++++++++++ src/sock_mgr.h | 6 ++++++ 5 files changed, 71 insertions(+), 3 deletions(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index 6e1e0e76c..dfd593757 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -70,10 +70,11 @@ void event_mgr::del_all_events(const std::string &tag) { int count = 0; for (const auto &event : this->event_map[tag]) { + int fd = event_get_fd(event); event_del(event); event_free(event); count++; - syslog(LOG_INFO, "event_mgr: Deleted event (fd=%d) of tag %s from %s", event_get_fd(event), tag.c_str(), this->name.c_str()); + syslog(LOG_INFO, "event_mgr: Deleted event (fd=%d) of tag %s from %s", fd, tag.c_str(), this->name.c_str()); } if (tag != "") { std::unordered_set &tagless_set = this->event_map[""]; @@ -88,6 +89,24 @@ void event_mgr::del_all_events(const std::string &tag) syslog(LOG_INFO, "event_mgr: Deleted %d events of tag %s for %s", count, tag.c_str(), this->name.c_str()); } +void event_mgr::suspend_all_events(const std::string &tag) +{ + for (const auto &event : this->event_map[tag]) { + event_del(event); + } +} + +int event_mgr::resume_all_events(const std::string &tag) +{ + for (const auto &event : this->event_map[tag]) { + if (event_add(event, NULL) < 0) { + this->suspend_all_events(tag); + return -1; + } + } + return 0; +} + /** * @code activate_all_events(tag, res); * diff --git a/src/event_mgr.h b/src/event_mgr.h index 90ff4a146..26a1b5135 100644 --- a/src/event_mgr.h +++ b/src/event_mgr.h @@ -12,6 +12,8 @@ class event_mgr { int init_base(); int add_event(struct event* event, const struct timeval *timeout, const std::string &tag=""); void del_all_events(const std::string &tag=""); + void suspend_all_events(const std::string &tag); + int resume_all_events(const std::string &tag); void activate_all_events(const std::string &tag="", int res=0); void free(); struct event_base* get_base(); diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index 7ec5d07a3..4353633dc 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -16,6 +16,8 @@ #include "dhcp_check_profile.h" /** to get dhcp/v6 check profile */ #include "util.h" +static constexpr int MAX_PACKETS_PER_CALLBACK = 64; + /** * @code _increase_cache_counter(ifname, sock, type); * @brief helper function to increase cache counter. Simple increase of counter, no complications. In the event of @@ -864,8 +866,12 @@ void callback_common(int fd, short event, void *arg) socklen_t slen = sizeof(sll); sock_info_t &sock_info = sock_mgr_get_sock_info(fd); - while ((buffer_sz = recvfrom(fd, sock_info.buffer, sock_info.snaplen, MSG_DONTWAIT, (struct sockaddr *)&sll, &slen)) > 0) - { + for (int packet_count = 0; packet_count < MAX_PACKETS_PER_CALLBACK; packet_count++) { + buffer_sz = recvfrom(fd, sock_info.buffer, sock_info.snaplen, MSG_DONTWAIT, + (struct sockaddr *)&sll, &slen); + if (buffer_sz <= 0) { + break; + } char ifname_buf[IF_NAMESIZE]; if (if_indextoname(sll.sll_ifindex, ifname_buf) == NULL) { syslog_debug(LOG_WARNING, "if_indextoname: invalid input interface index %d %s", sll.sll_ifindex, strerror(errno)); diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index 8d3e48d81..ef934fc44 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -36,6 +36,11 @@ static const char dhcpv6_outbound_filter[] = "outbound and ip6 and udp and (port /** Tags for different events, so we can triiger only one type */ static const char packet_handler_tag[] = "PacketHandler"; static const char cache_counter_updater_tag[] = "CacheCounterUpdater"; +static const char keepalive_tag[] = "Keepalive"; + +static void keepalive_callback(evutil_socket_t, short, void *) +{ +} /* sock fd to sock_info mapping */ std::unordered_map sock_map; @@ -385,6 +390,18 @@ int sock_mgr_init_event_mgr() sock_mgr_free_event_mgr(); return -1; } + struct event *keepalive_event = event_new(info.event_mgr_ptr->get_base(), -1, EV_PERSIST, + keepalive_callback, NULL); + struct timeval keepalive_interval = {.tv_sec = 3600, .tv_usec = 0}; + if (keepalive_event == NULL || + info.event_mgr_ptr->add_event(keepalive_event, &keepalive_interval, keepalive_tag) < 0) { + if (keepalive_event != NULL) { + event_free(keepalive_event); + } + syslog(LOG_ALERT, "Failed to initialize event manager keepalive %s", info.name); + sock_mgr_free_event_mgr(); + return -1; + } } return 0; @@ -432,6 +449,24 @@ void sock_mgr_unregister_packet_handler() } } +void sock_mgr_suspend_packet_handler() +{ + for (const auto &[sock, info] : sock_map) { + info.event_mgr_ptr->suspend_all_events(packet_handler_tag); + } +} + +int sock_mgr_resume_packet_handler() +{ + for (const auto &[sock, info] : sock_map) { + if (info.event_mgr_ptr->resume_all_events(packet_handler_tag) < 0) { + sock_mgr_suspend_packet_handler(); + return -1; + } + } + return 0; +} + int sock_mgr_register_cache_counter_updater(event_callback_fn callback) { syslog(LOG_INFO, "Registering cache counter updater for all sockets"); diff --git a/src/sock_mgr.h b/src/sock_mgr.h index 9619a9526..736ce9b8d 100644 --- a/src/sock_mgr.h +++ b/src/sock_mgr.h @@ -59,6 +59,12 @@ int sock_mgr_register_packet_handler(); /** Unregister packet handler for socket manager */ void sock_mgr_unregister_packet_handler(); +/** Temporarily suspend registered packet handlers */ +void sock_mgr_suspend_packet_handler(); + +/** Resume registered packet handlers */ +int sock_mgr_resume_packet_handler(); + /** Register cache counter updater callback for socket manager */ int sock_mgr_register_cache_counter_updater(event_callback_fn callback); From e3da99c92e14c5c8c245ea996b44a33702630bb8 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sat, 25 Jul 2026 23:06:44 +0000 Subject: [PATCH 02/15] [dhcpmon]: Harden packet event resume Reset sockaddr length for each batch receive and reject timeout events from the fd-only resume path. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 5 +++++ src/packet_handler.cpp | 1 + 2 files changed, 6 insertions(+) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index dfd593757..5d6e43f2f 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -99,6 +99,11 @@ void event_mgr::suspend_all_events(const std::string &tag) int event_mgr::resume_all_events(const std::string &tag) { for (const auto &event : this->event_map[tag]) { + if (event_get_fd(event) < 0) { + syslog(LOG_ALERT, "event_mgr: Cannot resume non-fd event with tag %s", tag.c_str()); + this->suspend_all_events(tag); + return -1; + } if (event_add(event, NULL) < 0) { this->suspend_all_events(tag); return -1; diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index 4353633dc..9f38ed5b2 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -867,6 +867,7 @@ void callback_common(int fd, short event, void *arg) sock_info_t &sock_info = sock_mgr_get_sock_info(fd); for (int packet_count = 0; packet_count < MAX_PACKETS_PER_CALLBACK; packet_count++) { + slen = sizeof(sll); buffer_sz = recvfrom(fd, sock_info.buffer, sock_info.snaplen, MSG_DONTWAIT, (struct sockaddr *)&sll, &slen); if (buffer_sz <= 0) { From 26953736102652f5559353a4c3c305f4dff41cc6 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 11:03:19 +1000 Subject: [PATCH 03/15] [dhcpmon]: Guard packet event suspension Reject untagged suspend/resume requests and log the event manager and fd when a packet event cannot be restored. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index 5d6e43f2f..b92cd919a 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -91,6 +91,11 @@ void event_mgr::del_all_events(const std::string &tag) void event_mgr::suspend_all_events(const std::string &tag) { + if (tag.empty()) { + syslog(LOG_ALERT, "event_mgr: Refusing to suspend untagged events for %s", + this->name.c_str()); + return; + } for (const auto &event : this->event_map[tag]) { event_del(event); } @@ -98,13 +103,21 @@ void event_mgr::suspend_all_events(const std::string &tag) int event_mgr::resume_all_events(const std::string &tag) { + if (tag.empty()) { + syslog(LOG_ALERT, "event_mgr: Refusing to resume untagged events for %s", + this->name.c_str()); + return -1; + } for (const auto &event : this->event_map[tag]) { if (event_get_fd(event) < 0) { - syslog(LOG_ALERT, "event_mgr: Cannot resume non-fd event with tag %s", tag.c_str()); + syslog(LOG_ALERT, "event_mgr: Cannot resume non-fd event with tag %s for %s", + tag.c_str(), this->name.c_str()); this->suspend_all_events(tag); return -1; } if (event_add(event, NULL) < 0) { + syslog(LOG_ALERT, "event_mgr: Failed to resume event (fd=%d) with tag %s for %s", + event_get_fd(event), tag.c_str(), this->name.c_str()); this->suspend_all_events(tag); return -1; } From a53f47c54eec0ca15f1dd6fe555ad3eff5b73f00 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 11:33:38 +1000 Subject: [PATCH 04/15] [dhcpmon]: Reject unknown packet event tags Look up tagged event sets without default insertion, logging unknown suspend tags and returning an error for unknown resume tags. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index b92cd919a..76f4b8b24 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -96,7 +96,13 @@ void event_mgr::suspend_all_events(const std::string &tag) this->name.c_str()); return; } - for (const auto &event : this->event_map[tag]) { + const auto tagged_events = this->event_map.find(tag); + if (tagged_events == this->event_map.end()) { + syslog(LOG_ALERT, "event_mgr: Cannot suspend unknown tag %s for %s", + tag.c_str(), this->name.c_str()); + return; + } + for (const auto &event : tagged_events->second) { event_del(event); } } @@ -108,7 +114,13 @@ int event_mgr::resume_all_events(const std::string &tag) this->name.c_str()); return -1; } - for (const auto &event : this->event_map[tag]) { + const auto tagged_events = this->event_map.find(tag); + if (tagged_events == this->event_map.end()) { + syslog(LOG_ALERT, "event_mgr: Cannot resume unknown tag %s for %s", + tag.c_str(), this->name.c_str()); + return -1; + } + for (const auto &event : tagged_events->second) { if (event_get_fd(event) < 0) { syslog(LOG_ALERT, "event_mgr: Cannot resume non-fd event with tag %s for %s", tag.c_str(), this->name.c_str()); From 9fb9cd676450659dee41f3c7608a7572f626247b Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 11:43:56 +1000 Subject: [PATCH 05/15] [dhcpmon]: Wait for in-flight packet callbacks Gate packet callbacks before deleting their events and take an exclusive quiesce lock so topology reconciliation cannot race a callback already executing on a socket event thread. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/packet_handler.cpp | 7 +++++++ src/sock_mgr.cpp | 37 ++++++++++++++++++++++++++++++------- src/sock_mgr.h | 6 ++++++ 3 files changed, 43 insertions(+), 7 deletions(-) diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index 9f38ed5b2..d52f2d048 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -861,6 +861,13 @@ void packet_handler_v6(int sock, const std::string &ifname, const dhcp_device_co void callback_common(int fd, short event, void *arg) { + if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + return; + } + std::shared_lock packet_handler_lock(packet_handler_quiesce_mutex); + if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + return; + } ssize_t buffer_sz; struct sockaddr_ll sll; socklen_t slen = sizeof(sll); diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index ef934fc44..c74dc01e9 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -11,6 +11,7 @@ #include #include #include +#include #include #include "sock_mgr.h" @@ -45,6 +46,10 @@ static void keepalive_callback(evutil_socket_t, short, void *) /* sock fd to sock_info mapping */ std::unordered_map sock_map; +std::shared_mutex packet_handler_quiesce_mutex; +std::atomic packet_handlers_enabled{true}; +static std::unique_lock packet_handler_quiesce_lock; + extern std::shared_ptr mCountersDbPtr; extern std::string downstream_ifname; @@ -451,20 +456,38 @@ void sock_mgr_unregister_packet_handler() void sock_mgr_suspend_packet_handler() { - for (const auto &[sock, info] : sock_map) { - info.event_mgr_ptr->suspend_all_events(packet_handler_tag); + if (packet_handler_quiesce_lock.owns_lock()) { + syslog(LOG_ALERT, "Packet handlers are already suspended"); + return; + } + packet_handlers_enabled.store(false, std::memory_order_release); + for (const auto &entry : sock_map) { + entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); } + packet_handler_quiesce_lock = std::unique_lock(packet_handler_quiesce_mutex); } int sock_mgr_resume_packet_handler() { - for (const auto &[sock, info] : sock_map) { - if (info.event_mgr_ptr->resume_all_events(packet_handler_tag) < 0) { - sock_mgr_suspend_packet_handler(); - return -1; + if (!packet_handler_quiesce_lock.owns_lock()) { + syslog(LOG_ALERT, "Packet handlers are not suspended"); + return -1; + } + int result = 0; + for (const auto &entry : sock_map) { + if (entry.second.event_mgr_ptr->resume_all_events(packet_handler_tag) < 0) { + for (const auto &suspended_entry : sock_map) { + suspended_entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); + } + result = -1; + break; } } - return 0; + if (result == 0) { + packet_handlers_enabled.store(true, std::memory_order_release); + } + packet_handler_quiesce_lock.unlock(); + return result; } int sock_mgr_register_cache_counter_updater(event_callback_fn callback) diff --git a/src/sock_mgr.h b/src/sock_mgr.h index 736ce9b8d..e41eccf57 100644 --- a/src/sock_mgr.h +++ b/src/sock_mgr.h @@ -9,7 +9,9 @@ #ifndef SOCKET_MANAGER_H_ #define SOCKET_MANAGER_H_ +#include #include +#include #include #include #include @@ -41,6 +43,10 @@ typedef struct { /** sock file descriptors, serve as the identifier of all related information described in sock_info_t */ extern int rx_sock, tx_sock, rx_sock_v6, tx_sock_v6; +/** Guards in-flight packet callbacks while topology and counters are reconciled */ +extern std::shared_mutex packet_handler_quiesce_mutex; +extern std::atomic packet_handlers_enabled; + /** Initialize socket manager with given snaplen */ int sock_mgr_init(uint32_t snaplen); From 7175f81ae109a09a1d210e4bd3281d84ce5e3108 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:03:35 +1000 Subject: [PATCH 06/15] [dhcpmon]: Avoid implicit event tag creation Use explicit tag lookups for event deletion and activation so unknown tags are logged without mutating the event map or reporting a silent no-op. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 31 +++++++++++++++++++++++-------- 1 file changed, 23 insertions(+), 8 deletions(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index 76f4b8b24..84da07fd5 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -69,19 +69,26 @@ int event_mgr::add_event(struct event* event, const struct timeval *timeout, con void event_mgr::del_all_events(const std::string &tag) { int count = 0; - for (const auto &event : this->event_map[tag]) { + const auto tagged_events = this->event_map.find(tag); + if (tagged_events == this->event_map.end()) { + if (!tag.empty()) { + syslog(LOG_WARNING, "event_mgr: Cannot delete unknown tag %s for %s", + tag.c_str(), this->name.c_str()); + } + return; + } + auto all_events = this->event_map.find(""); + for (const auto &event : tagged_events->second) { int fd = event_get_fd(event); + if (!tag.empty() && all_events != this->event_map.end()) { + all_events->second.erase(event); + } event_del(event); event_free(event); count++; syslog(LOG_INFO, "event_mgr: Deleted event (fd=%d) of tag %s from %s", fd, tag.c_str(), this->name.c_str()); } - if (tag != "") { - std::unordered_set &tagless_set = this->event_map[""]; - std::unordered_set &tagged_set = this->event_map[tag]; - for (const auto &event : tagged_set) { - tagless_set.erase(event); - } + if (!tag.empty()) { this->event_map.erase(tag); } else { this->event_map.clear(); @@ -146,7 +153,15 @@ int event_mgr::resume_all_events(const std::string &tag) */ void event_mgr::activate_all_events(const std::string &tag, int res) { - for (const auto &event : this->event_map[tag]) { + const auto tagged_events = this->event_map.find(tag); + if (tagged_events == this->event_map.end()) { + if (!tag.empty()) { + syslog(LOG_WARNING, "event_mgr: Cannot activate unknown tag %s for %s", + tag.c_str(), this->name.c_str()); + } + return; + } + for (const auto &event : tagged_events->second) { event_active(event, res, 0); syslog(LOG_INFO, "event_mgr: Activated event (fd=%d) of tag %s from %s", event_get_fd(event), tag.c_str(), this->name.c_str()); } From 9d0ee6a2c7f19e237504bdd292fdd3f7f0601826 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:17:01 +1000 Subject: [PATCH 07/15] [dhcpmon]: Avoid blocking callbacks during quiesce Use a non-blocking shared-lock attempt so an event-loop callback exits immediately when topology reconciliation owns the exclusive quiesce lock. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/packet_handler.cpp | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index d52f2d048..742b388fd 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -8,6 +8,7 @@ #include #include #include +#include #include "packet_handler.h" @@ -864,8 +865,10 @@ void callback_common(int fd, short event, void *arg) if (!packet_handlers_enabled.load(std::memory_order_acquire)) { return; } - std::shared_lock packet_handler_lock(packet_handler_quiesce_mutex); - if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + std::shared_lock packet_handler_lock(packet_handler_quiesce_mutex, + std::try_to_lock); + if (!packet_handler_lock.owns_lock() || + !packet_handlers_enabled.load(std::memory_order_acquire)) { return; } ssize_t buffer_sz; From 55a9104b5259b9a7ea35d2cded4a48ee5a7f25af Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:17:01 +1000 Subject: [PATCH 08/15] [dhcpmon]: Keep event cleanup idempotent Treat deletion of an absent event tag as the expected cleanup no-op while retaining explicit lookup and avoiding map mutation. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 4 ---- 1 file changed, 4 deletions(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index 84da07fd5..a794cee1a 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -71,10 +71,6 @@ void event_mgr::del_all_events(const std::string &tag) int count = 0; const auto tagged_events = this->event_map.find(tag); if (tagged_events == this->event_map.end()) { - if (!tag.empty()) { - syslog(LOG_WARNING, "event_mgr: Cannot delete unknown tag %s for %s", - tag.c_str(), this->name.c_str()); - } return; } auto all_events = this->event_map.find(""); From ed1661f9e0195900fb4be01b48355dddb12a14ee Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:29:00 +1000 Subject: [PATCH 09/15] [dhcpmon]: Synchronize counter sampling and updates Take the exclusive counter-state lock for health snapshots, DB synchronization, clear-counter cache updates, and signal status reads so packet callbacks cannot race those operations. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/dhcp_mon.cpp | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/src/dhcp_mon.cpp b/src/dhcp_mon.cpp index ef4f6623d..42551201d 100644 --- a/src/dhcp_mon.cpp +++ b/src/dhcp_mon.cpp @@ -212,9 +212,12 @@ static void cleanup_stale_db_counters() static void signal_callback(evutil_socket_t fd, short event, void *arg) { syslog(LOG_INFO, "Received signal: %s", strsignal(fd)); - - dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT); - dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT_V6); + + { + std::unique_lock counter_lock(packet_handler_quiesce_mutex); + dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT); + dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT_V6); + } if ((fd == SIGTERM) || (fd == SIGINT)) { syslog(LOG_INFO, "Received signal to stop dhcpmon"); @@ -223,6 +226,7 @@ static void signal_callback(evutil_socket_t fd, short event, void *arg) if (fd == SIGUSR1) { // we need to sync cache counter from COUNTERS_DB syslog(LOG_INFO, "Received signal to stop writing to DB counter"); + std::unique_lock counter_lock(packet_handler_quiesce_mutex); std::lock_guard lock(db_sync_mutex); sock_mgr_pause_write_cache_to_db(); syslog(LOG_INFO, "Stopped writing to DB counter"); @@ -260,6 +264,7 @@ static void update_cache_counter_callback(evutil_socket_t fd, short event, void syslog(LOG_INFO, "Start updating %s cache counter from DB counter", sock_info.name); + std::unique_lock counter_lock(packet_handler_quiesce_mutex); std::lock_guard lock(db_sync_mutex); // can only sync db to cache counter and db updater is paused, otherwise its unexpected @@ -392,6 +397,7 @@ static void update_cache_counter_callback(evutil_socket_t fd, short event, void static void timeout_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received timeout signal for DHCP relay health check"); + std::unique_lock counter_lock(packet_handler_quiesce_mutex); dhcp_devman_print_all_status_debug(DHCP_COUNTERS_CURRENT); dhcp_devman_print_all_status_debug(DHCP_COUNTERS_SNAPSHOT); @@ -418,6 +424,7 @@ static void db_update_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received db update signal"); syslog_debug(LOG_INFO, "Sync cache counter to DB counter"); + std::unique_lock counter_lock(packet_handler_quiesce_mutex); std::lock_guard lock(db_sync_mutex); // If there is clear counter going on and its been longer than expected // consider the clear counter operation failed so we don't block db update forever From d87d651edf3adab3367bb1c303adc5d483934814 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:37:21 +1000 Subject: [PATCH 10/15] [dhcpmon]: Guarantee counter writer progress Announce pending exclusive counter-state operations before locking so packet callbacks stop admitting shared readers and cannot starve health, DB, clear, or topology writers. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/dhcp_mon.cpp | 10 +++++----- src/packet_handler.cpp | 6 ++++-- src/sock_mgr.cpp | 18 ++++++++++++++++++ src/sock_mgr.h | 14 ++++++++++++++ 4 files changed, 41 insertions(+), 7 deletions(-) diff --git a/src/dhcp_mon.cpp b/src/dhcp_mon.cpp index 42551201d..143154b0a 100644 --- a/src/dhcp_mon.cpp +++ b/src/dhcp_mon.cpp @@ -214,7 +214,7 @@ static void signal_callback(evutil_socket_t fd, short event, void *arg) syslog(LOG_INFO, "Received signal: %s", strsignal(fd)); { - std::unique_lock counter_lock(packet_handler_quiesce_mutex); + counter_state_write_lock counter_lock; dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT); dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT_V6); } @@ -226,7 +226,7 @@ static void signal_callback(evutil_socket_t fd, short event, void *arg) if (fd == SIGUSR1) { // we need to sync cache counter from COUNTERS_DB syslog(LOG_INFO, "Received signal to stop writing to DB counter"); - std::unique_lock counter_lock(packet_handler_quiesce_mutex); + counter_state_write_lock counter_lock; std::lock_guard lock(db_sync_mutex); sock_mgr_pause_write_cache_to_db(); syslog(LOG_INFO, "Stopped writing to DB counter"); @@ -264,7 +264,7 @@ static void update_cache_counter_callback(evutil_socket_t fd, short event, void syslog(LOG_INFO, "Start updating %s cache counter from DB counter", sock_info.name); - std::unique_lock counter_lock(packet_handler_quiesce_mutex); + counter_state_write_lock counter_lock; std::lock_guard lock(db_sync_mutex); // can only sync db to cache counter and db updater is paused, otherwise its unexpected @@ -397,7 +397,7 @@ static void update_cache_counter_callback(evutil_socket_t fd, short event, void static void timeout_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received timeout signal for DHCP relay health check"); - std::unique_lock counter_lock(packet_handler_quiesce_mutex); + counter_state_write_lock counter_lock; dhcp_devman_print_all_status_debug(DHCP_COUNTERS_CURRENT); dhcp_devman_print_all_status_debug(DHCP_COUNTERS_SNAPSHOT); @@ -424,7 +424,7 @@ static void db_update_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received db update signal"); syslog_debug(LOG_INFO, "Sync cache counter to DB counter"); - std::unique_lock counter_lock(packet_handler_quiesce_mutex); + counter_state_write_lock counter_lock; std::lock_guard lock(db_sync_mutex); // If there is clear counter going on and its been longer than expected // consider the clear counter operation failed so we don't block db update forever diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index 742b388fd..5cbd12363 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -862,13 +862,15 @@ void packet_handler_v6(int sock, const std::string &ifname, const dhcp_device_co void callback_common(int fd, short event, void *arg) { - if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + if (!packet_handlers_enabled.load(std::memory_order_acquire) || + counter_state_writers_pending.load(std::memory_order_acquire) > 0) { return; } std::shared_lock packet_handler_lock(packet_handler_quiesce_mutex, std::try_to_lock); if (!packet_handler_lock.owns_lock() || - !packet_handlers_enabled.load(std::memory_order_acquire)) { + !packet_handlers_enabled.load(std::memory_order_acquire) || + counter_state_writers_pending.load(std::memory_order_acquire) > 0) { return; } ssize_t buffer_sz; diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index c74dc01e9..71b85485e 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -48,12 +48,30 @@ std::unordered_map sock_map; std::shared_mutex packet_handler_quiesce_mutex; std::atomic packet_handlers_enabled{true}; +std::atomic counter_state_writers_pending{0}; static std::unique_lock packet_handler_quiesce_lock; extern std::shared_ptr mCountersDbPtr; extern std::string downstream_ifname; +counter_state_write_lock::counter_state_write_lock() +{ + counter_state_writers_pending.fetch_add(1, std::memory_order_acq_rel); + try { + lock = std::unique_lock(packet_handler_quiesce_mutex); + } catch (...) { + counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel); + throw; + } +} + +counter_state_write_lock::~counter_state_write_lock() +{ + lock.unlock(); + counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel); +} + /** * @code opensocket(); * diff --git a/src/sock_mgr.h b/src/sock_mgr.h index e41eccf57..52b265526 100644 --- a/src/sock_mgr.h +++ b/src/sock_mgr.h @@ -10,6 +10,7 @@ #define SOCKET_MANAGER_H_ #include +#include #include #include #include @@ -46,6 +47,19 @@ extern int rx_sock, tx_sock, rx_sock_v6, tx_sock_v6; /** Guards in-flight packet callbacks while topology and counters are reconciled */ extern std::shared_mutex packet_handler_quiesce_mutex; extern std::atomic packet_handlers_enabled; +extern std::atomic counter_state_writers_pending; + +class counter_state_write_lock +{ + public: + counter_state_write_lock(); + ~counter_state_write_lock(); + counter_state_write_lock(const counter_state_write_lock &) = delete; + counter_state_write_lock &operator=(const counter_state_write_lock &) = delete; + + private: + std::unique_lock lock; +}; /** Initialize socket manager with given snaplen */ int sock_mgr_init(uint32_t snaplen); From b340941baa73184fdabd8337a8efa1cbc6433901 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 12:55:31 +1000 Subject: [PATCH 11/15] [dhcpmon]: Block packet callbacks without spinning Use a condition-backed reader gate for pending counter writers, make topology quiesce transitions explicit and recoverable, and surface counter-lock failures without unwinding through libevent callbacks. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/dhcp_mon.cpp | 18 +++++- src/packet_handler.cpp | 12 +--- src/sock_mgr.cpp | 130 +++++++++++++++++++++++++++++++++++------ src/sock_mgr.h | 13 ++++- 4 files changed, 143 insertions(+), 30 deletions(-) diff --git a/src/dhcp_mon.cpp b/src/dhcp_mon.cpp index 143154b0a..3504737a8 100644 --- a/src/dhcp_mon.cpp +++ b/src/dhcp_mon.cpp @@ -215,8 +215,10 @@ static void signal_callback(evutil_socket_t fd, short event, void *arg) { counter_state_write_lock counter_lock; - dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT); - dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT_V6); + if (counter_lock.owns_lock()) { + dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT); + dhcp_devman_print_all_status(DHCP_COUNTERS_CURRENT_V6); + } } if ((fd == SIGTERM) || (fd == SIGINT)) { @@ -227,6 +229,9 @@ static void signal_callback(evutil_socket_t fd, short event, void *arg) // we need to sync cache counter from COUNTERS_DB syslog(LOG_INFO, "Received signal to stop writing to DB counter"); counter_state_write_lock counter_lock; + if (!counter_lock.owns_lock()) { + return; + } std::lock_guard lock(db_sync_mutex); sock_mgr_pause_write_cache_to_db(); syslog(LOG_INFO, "Stopped writing to DB counter"); @@ -265,6 +270,9 @@ static void update_cache_counter_callback(evutil_socket_t fd, short event, void syslog(LOG_INFO, "Start updating %s cache counter from DB counter", sock_info.name); counter_state_write_lock counter_lock; + if (!counter_lock.owns_lock()) { + return; + } std::lock_guard lock(db_sync_mutex); // can only sync db to cache counter and db updater is paused, otherwise its unexpected @@ -398,6 +406,9 @@ static void timeout_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received timeout signal for DHCP relay health check"); counter_state_write_lock counter_lock; + if (!counter_lock.owns_lock()) { + return; + } dhcp_devman_print_all_status_debug(DHCP_COUNTERS_CURRENT); dhcp_devman_print_all_status_debug(DHCP_COUNTERS_SNAPSHOT); @@ -425,6 +436,9 @@ static void db_update_callback(evutil_socket_t fd, short event, void *arg) syslog_debug(LOG_INFO, "Received db update signal"); syslog_debug(LOG_INFO, "Sync cache counter to DB counter"); counter_state_write_lock counter_lock; + if (!counter_lock.owns_lock()) { + return; + } std::lock_guard lock(db_sync_mutex); // If there is clear counter going on and its been longer than expected // consider the clear counter operation failed so we don't block db update forever diff --git a/src/packet_handler.cpp b/src/packet_handler.cpp index 5cbd12363..aa564c1a0 100644 --- a/src/packet_handler.cpp +++ b/src/packet_handler.cpp @@ -8,7 +8,6 @@ #include #include #include -#include #include "packet_handler.h" @@ -862,15 +861,8 @@ void packet_handler_v6(int sock, const std::string &ifname, const dhcp_device_co void callback_common(int fd, short event, void *arg) { - if (!packet_handlers_enabled.load(std::memory_order_acquire) || - counter_state_writers_pending.load(std::memory_order_acquire) > 0) { - return; - } - std::shared_lock packet_handler_lock(packet_handler_quiesce_mutex, - std::try_to_lock); - if (!packet_handler_lock.owns_lock() || - !packet_handlers_enabled.load(std::memory_order_acquire) || - counter_state_writers_pending.load(std::memory_order_acquire) > 0) { + counter_state_read_lock counter_lock; + if (!counter_lock.owns_lock()) { return; } ssize_t buffer_sz; diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index 71b85485e..8042c7e2e 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -11,7 +11,9 @@ #include #include #include +#include #include +#include #include #include "sock_mgr.h" @@ -50,26 +52,97 @@ std::shared_mutex packet_handler_quiesce_mutex; std::atomic packet_handlers_enabled{true}; std::atomic counter_state_writers_pending{0}; static std::unique_lock packet_handler_quiesce_lock; +static std::mutex counter_state_wait_mutex; +static std::condition_variable counter_state_wait_cv; extern std::shared_ptr mCountersDbPtr; extern std::string downstream_ifname; +static void set_packet_handlers_enabled(bool enabled) +{ + { + std::lock_guard wait_lock(counter_state_wait_mutex); + packet_handlers_enabled.store(enabled, std::memory_order_release); + } + counter_state_wait_cv.notify_all(); +} + counter_state_write_lock::counter_state_write_lock() { - counter_state_writers_pending.fetch_add(1, std::memory_order_acq_rel); + { + std::lock_guard wait_lock(counter_state_wait_mutex); + counter_state_writers_pending.fetch_add(1, std::memory_order_acq_rel); + } try { lock = std::unique_lock(packet_handler_quiesce_mutex); - } catch (...) { - counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel); - throw; + } catch (const std::system_error &e) { + bool notify = false; + { + std::lock_guard wait_lock(counter_state_wait_mutex); + notify = counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel) == 1; + } + if (notify) { + counter_state_wait_cv.notify_all(); + } + syslog(LOG_ALERT, "Failed to lock DHCP counter state: %s", e.what()); } } counter_state_write_lock::~counter_state_write_lock() { + if (!lock.owns_lock()) { + return; + } lock.unlock(); - counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel); + bool notify = false; + { + std::lock_guard wait_lock(counter_state_wait_mutex); + notify = counter_state_writers_pending.fetch_sub(1, std::memory_order_acq_rel) == 1; + } + if (notify) { + counter_state_wait_cv.notify_all(); + } +} + +bool counter_state_write_lock::owns_lock() const +{ + return lock.owns_lock(); +} + +counter_state_read_lock::counter_state_read_lock() +{ + while (packet_handlers_enabled.load(std::memory_order_acquire)) { + { + std::unique_lock wait_lock(counter_state_wait_mutex); + counter_state_wait_cv.wait(wait_lock, [] { + return !packet_handlers_enabled.load(std::memory_order_acquire) || + counter_state_writers_pending.load(std::memory_order_acquire) == 0; + }); + } + if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + return; + } + try { + lock = std::shared_lock(packet_handler_quiesce_mutex); + } catch (const std::system_error &e) { + syslog(LOG_ALERT, "Failed to lock DHCP counter state for packet handling: %s", e.what()); + return; + } + if (!packet_handlers_enabled.load(std::memory_order_acquire)) { + lock.unlock(); + return; + } + if (counter_state_writers_pending.load(std::memory_order_acquire) == 0) { + return; + } + lock.unlock(); + } +} + +bool counter_state_read_lock::owns_lock() const +{ + return lock.owns_lock(); } /** @@ -472,17 +545,37 @@ void sock_mgr_unregister_packet_handler() } } -void sock_mgr_suspend_packet_handler() +int sock_mgr_suspend_packet_handler() { if (packet_handler_quiesce_lock.owns_lock()) { syslog(LOG_ALERT, "Packet handlers are already suspended"); - return; + return -1; } - packet_handlers_enabled.store(false, std::memory_order_release); for (const auto &entry : sock_map) { entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); } - packet_handler_quiesce_lock = std::unique_lock(packet_handler_quiesce_mutex); + set_packet_handlers_enabled(false); + try { + packet_handler_quiesce_lock = std::unique_lock(packet_handler_quiesce_mutex); + } catch (const std::system_error &e) { + syslog(LOG_ALERT, "Failed to quiesce packet handlers: %s", e.what()); + set_packet_handlers_enabled(true); + int restore_result = 0; + for (const auto &entry : sock_map) { + if (entry.second.event_mgr_ptr->resume_all_events(packet_handler_tag) < 0) { + restore_result = -1; + } + } + if (restore_result < 0) { + set_packet_handlers_enabled(false); + for (const auto &entry : sock_map) { + entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); + } + syslog(LOG_ALERT, "Failed to restore packet handlers after quiesce failure"); + } + return -1; + } + return 0; } int sock_mgr_resume_packet_handler() @@ -491,21 +584,24 @@ int sock_mgr_resume_packet_handler() syslog(LOG_ALERT, "Packet handlers are not suspended"); return -1; } - int result = 0; + set_packet_handlers_enabled(true); + packet_handler_quiesce_lock.unlock(); + for (const auto &entry : sock_map) { if (entry.second.event_mgr_ptr->resume_all_events(packet_handler_tag) < 0) { + set_packet_handlers_enabled(false); for (const auto &suspended_entry : sock_map) { suspended_entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); } - result = -1; - break; + try { + packet_handler_quiesce_lock = std::unique_lock(packet_handler_quiesce_mutex); + } catch (const std::system_error &e) { + syslog(LOG_ALERT, "Failed to restore packet quiesce lock after resume failure: %s", e.what()); + } + return -1; } } - if (result == 0) { - packet_handlers_enabled.store(true, std::memory_order_release); - } - packet_handler_quiesce_lock.unlock(); - return result; + return 0; } int sock_mgr_register_cache_counter_updater(event_callback_fn callback) diff --git a/src/sock_mgr.h b/src/sock_mgr.h index 52b265526..7801d70f2 100644 --- a/src/sock_mgr.h +++ b/src/sock_mgr.h @@ -54,6 +54,7 @@ class counter_state_write_lock public: counter_state_write_lock(); ~counter_state_write_lock(); + bool owns_lock() const; counter_state_write_lock(const counter_state_write_lock &) = delete; counter_state_write_lock &operator=(const counter_state_write_lock &) = delete; @@ -61,6 +62,16 @@ class counter_state_write_lock std::unique_lock lock; }; +class counter_state_read_lock +{ + public: + counter_state_read_lock(); + bool owns_lock() const; + + private: + std::shared_lock lock; +}; + /** Initialize socket manager with given snaplen */ int sock_mgr_init(uint32_t snaplen); @@ -80,7 +91,7 @@ int sock_mgr_register_packet_handler(); void sock_mgr_unregister_packet_handler(); /** Temporarily suspend registered packet handlers */ -void sock_mgr_suspend_packet_handler(); +int sock_mgr_suspend_packet_handler(); /** Resume registered packet handlers */ int sock_mgr_resume_packet_handler(); From 626180a4f64b3e541201c8d876b63206c60042e4 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 13:21:22 +1000 Subject: [PATCH 12/15] [dhcpmon]: Keep packet callback fast path lock-free Try the shared counter lock directly when no writer is pending, using the condition-variable wait path only when quiescing or writer contention requires it. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/sock_mgr.cpp | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index 8042c7e2e..b316d963d 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -112,6 +112,25 @@ bool counter_state_write_lock::owns_lock() const counter_state_read_lock::counter_state_read_lock() { + if (packet_handlers_enabled.load(std::memory_order_acquire) && + counter_state_writers_pending.load(std::memory_order_acquire) == 0) { + try { + lock = std::shared_lock(packet_handler_quiesce_mutex, + std::try_to_lock); + } catch (const std::system_error &e) { + syslog(LOG_ALERT, "Failed to lock DHCP counter state for packet handling: %s", e.what()); + return; + } + if (lock.owns_lock() && + packet_handlers_enabled.load(std::memory_order_acquire) && + counter_state_writers_pending.load(std::memory_order_acquire) == 0) { + return; + } + if (lock.owns_lock()) { + lock.unlock(); + } + } + while (packet_handlers_enabled.load(std::memory_order_acquire)) { { std::unique_lock wait_lock(counter_state_wait_mutex); From b388c3620a9ce20bf92357c098f94b95c8c05628 Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 13:52:11 +1000 Subject: [PATCH 13/15] [dhcpmon]: Write DB counters from a stable snapshot Copy counter maps while holding counter-state and DB synchronization locks, then release packet callbacks before Redis I/O while retaining DB serialization through writeback. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/dhcp_mon.cpp | 37 +++++++++++++++++++++---------------- src/sock_mgr.cpp | 21 ++++++++++++++++++--- src/sock_mgr.h | 5 +++++ 3 files changed, 44 insertions(+), 19 deletions(-) diff --git a/src/dhcp_mon.cpp b/src/dhcp_mon.cpp index 3504737a8..10bf9c561 100644 --- a/src/dhcp_mon.cpp +++ b/src/dhcp_mon.cpp @@ -435,26 +435,31 @@ static void db_update_callback(evutil_socket_t fd, short event, void *arg) { syslog_debug(LOG_INFO, "Received db update signal"); syslog_debug(LOG_INFO, "Sync cache counter to DB counter"); - counter_state_write_lock counter_lock; - if (!counter_lock.owns_lock()) { - return; - } - std::lock_guard lock(db_sync_mutex); - // If there is clear counter going on and its been longer than expected - // consider the clear counter operation failed so we don't block db update forever - if (!sock_mgr_pause_write_cache_to_db_all_cleared() && last_update_time != default_time_point) { - auto now = std::chrono::steady_clock::now(); - auto elapsed = std::chrono::duration_cast(now - last_update_time); - if (elapsed.count() >= clear_counter_timeout) { - syslog(LOG_WARNING, "Clear counter going on for too long, abort clear counter"); - sock_mgr_clear_pause_write_cache_to_db(); - } else { - syslog(LOG_INFO, "Clear counter is ongoing, skip syncing write cache counter to DB counter"); + socket_counters_t counters_by_socket; + std::unique_lock lock; + { + counter_state_write_lock counter_lock; + if (!counter_lock.owns_lock()) { return; } + lock = std::unique_lock(db_sync_mutex); + // If there is clear counter going on and its been longer than expected + // consider the clear counter operation failed so we don't block db update forever + if (!sock_mgr_pause_write_cache_to_db_all_cleared() && last_update_time != default_time_point) { + auto now = std::chrono::steady_clock::now(); + auto elapsed = std::chrono::duration_cast(now - last_update_time); + if (elapsed.count() >= clear_counter_timeout) { + syslog(LOG_WARNING, "Clear counter going on for too long, abort clear counter"); + sock_mgr_clear_pause_write_cache_to_db(); + } else { + syslog(LOG_INFO, "Clear counter is ongoing, skip syncing write cache counter to DB counter"); + return; + } + } + counters_by_socket = sock_mgr_copy_cache_counters(); } last_update_time = std::chrono::steady_clock::now(); - sock_mgr_update_db_counters(); + sock_mgr_update_db_counters(counters_by_socket); cleanup_stale_db_counters(); syslog_debug(LOG_INFO, "Successfully synced cache counter to DB counter"); } diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index b316d963d..33f0cb2a8 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -805,17 +805,27 @@ bool sock_mgr_all_cache_counters_initialized(const std::string &ifname) return true; } -void sock_mgr_update_db_counters() +socket_counters_t sock_mgr_copy_cache_counters() +{ + socket_counters_t counters_by_socket; + for (const auto &[sock, info] : sock_map) { + counters_by_socket.emplace(sock, info.all_counters); + } + return counters_by_socket; +} + +void sock_mgr_update_db_counters(const socket_counters_t &counters_by_socket) { syslog_debug(LOG_INFO, "Updating all cache counters to DB counters"); - for (const auto &[sock, info] : sock_map) { + for (const auto &[sock, all_counters] : counters_by_socket) { + const sock_info_t &info = sock_mgr_get_sock_info(sock); syslog_debug(LOG_INFO, "Start updating socket %d %s DB counter from cache counter", sock, info.name); int msg_type_count = info.is_v6 ? DHCPV6_MESSAGE_TYPE_COUNT : DHCP_MESSAGE_TYPE_COUNT; const std::string *msg_type_name = info.is_v6 ? db_counter_name_v6 : db_counter_name; std::string all_ifname; std::string all_skipped_ifname; - for (const auto &[ifname, counter] : info.all_counters) { + for (const auto &[ifname, counter] : all_counters) { if (is_agg_counter(ifname) == true) { all_skipped_ifname += ifname + ", "; continue; @@ -830,4 +840,9 @@ void sock_mgr_update_db_counters() syslog_debug(LOG_INFO, "Skipped aggregated device counter entry of %sfor downstream vlan %s", all_skipped_ifname.c_str(), downstream_ifname.c_str()); } +} + +void sock_mgr_update_db_counters() +{ + sock_mgr_update_db_counters(sock_mgr_copy_cache_counters()); } \ No newline at end of file diff --git a/src/sock_mgr.h b/src/sock_mgr.h index 7801d70f2..3fd882159 100644 --- a/src/sock_mgr.h +++ b/src/sock_mgr.h @@ -22,6 +22,7 @@ typedef std::unordered_map counter_t; typedef std::unordered_map all_counters_t; +typedef std::unordered_map socket_counters_t; /** struct for socket information */ typedef struct { @@ -143,5 +144,9 @@ bool sock_mgr_all_cache_counters_initialized(const std::string &ifname); /** Update database counters from cache counters for all sockets */ void sock_mgr_update_db_counters(); +void sock_mgr_update_db_counters(const socket_counters_t &counters_by_socket); + +/** Copy cache counters for all sockets */ +socket_counters_t sock_mgr_copy_cache_counters(); #endif /* SOCKET_MANAGER_H_ */ From c576bf90961745b1f862d7702d5f56d60ed2e52b Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 14:14:15 +1000 Subject: [PATCH 14/15] [dhcpmon]: Skip unknown counter snapshot sockets Validate snapshot socket keys before looking up metadata so malformed or stale caller data cannot throw from sock_map.at(). Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/sock_mgr.cpp | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index 33f0cb2a8..ef758eefc 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -819,7 +819,12 @@ void sock_mgr_update_db_counters(const socket_counters_t &counters_by_socket) syslog_debug(LOG_INFO, "Updating all cache counters to DB counters"); for (const auto &[sock, all_counters] : counters_by_socket) { - const sock_info_t &info = sock_mgr_get_sock_info(sock); + const auto sock_info = sock_map.find(sock); + if (sock_info == sock_map.end()) { + syslog(LOG_WARNING, "Skip DB counter snapshot for unknown socket %d", sock); + continue; + } + const sock_info_t &info = sock_info->second; syslog_debug(LOG_INFO, "Start updating socket %d %s DB counter from cache counter", sock, info.name); int msg_type_count = info.is_v6 ? DHCPV6_MESSAGE_TYPE_COUNT : DHCP_MESSAGE_TYPE_COUNT; const std::string *msg_type_name = info.is_v6 ? db_counter_name_v6 : db_counter_name; From c447b5dfa9a28c515bc20eff07a45c8355b3c9fd Mon Sep 17 00:00:00 2001 From: Xichen96 Date: Sun, 26 Jul 2026 15:51:24 +1000 Subject: [PATCH 15/15] [dhcpmon]: Propagate packet event suspend failures Reject non-resumable event tags before deletion, roll back partial event_del failures, and surface suspend errors through the socket manager for coherent fail-stop recovery. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 39f979be-d826-4d5c-949a-f20abb58bb83 Signed-off-by: Xichen96 --- src/event_mgr.cpp | 29 +++++++++++++++++++++++++---- src/event_mgr.h | 2 +- src/sock_mgr.cpp | 34 ++++++++++++++++++++++++++++++---- 3 files changed, 56 insertions(+), 9 deletions(-) diff --git a/src/event_mgr.cpp b/src/event_mgr.cpp index a794cee1a..575b3ac15 100644 --- a/src/event_mgr.cpp +++ b/src/event_mgr.cpp @@ -1,4 +1,5 @@ #include +#include #include "event_mgr.h" @@ -92,22 +93,42 @@ void event_mgr::del_all_events(const std::string &tag) syslog(LOG_INFO, "event_mgr: Deleted %d events of tag %s for %s", count, tag.c_str(), this->name.c_str()); } -void event_mgr::suspend_all_events(const std::string &tag) +int event_mgr::suspend_all_events(const std::string &tag) { if (tag.empty()) { syslog(LOG_ALERT, "event_mgr: Refusing to suspend untagged events for %s", this->name.c_str()); - return; + return -1; } const auto tagged_events = this->event_map.find(tag); if (tagged_events == this->event_map.end()) { syslog(LOG_ALERT, "event_mgr: Cannot suspend unknown tag %s for %s", tag.c_str(), this->name.c_str()); - return; + return -1; } for (const auto &event : tagged_events->second) { - event_del(event); + if (event_get_fd(event) < 0) { + syslog(LOG_ALERT, "event_mgr: Cannot suspend non-fd event with tag %s for %s", + tag.c_str(), this->name.c_str()); + return -1; + } } + std::vector deleted_events; + for (const auto &event : tagged_events->second) { + if (event_del(event) < 0) { + bool restore_failed = false; + for (struct event *deleted_event : deleted_events) { + if (event_add(deleted_event, NULL) < 0) { + restore_failed = true; + } + } + syslog(LOG_ALERT, "event_mgr: Failed to suspend event (fd=%d) with tag %s for %s", + event_get_fd(event), tag.c_str(), this->name.c_str()); + return restore_failed ? -2 : -1; + } + deleted_events.push_back(event); + } + return 0; } int event_mgr::resume_all_events(const std::string &tag) diff --git a/src/event_mgr.h b/src/event_mgr.h index 26a1b5135..cca15edf8 100644 --- a/src/event_mgr.h +++ b/src/event_mgr.h @@ -12,7 +12,7 @@ class event_mgr { int init_base(); int add_event(struct event* event, const struct timeval *timeout, const std::string &tag=""); void del_all_events(const std::string &tag=""); - void suspend_all_events(const std::string &tag); + int suspend_all_events(const std::string &tag); int resume_all_events(const std::string &tag); void activate_all_events(const std::string &tag="", int res=0); void free(); diff --git a/src/sock_mgr.cpp b/src/sock_mgr.cpp index ef758eefc..e20f368d8 100644 --- a/src/sock_mgr.cpp +++ b/src/sock_mgr.cpp @@ -15,6 +15,7 @@ #include #include #include +#include #include "sock_mgr.h" @@ -36,7 +37,7 @@ static const char dhcp_outbound_filter[] = "outbound and ip and udp and (port 67 static const char dhcpv6_inbound_filter[] = "inbound and ip6 and udp and (port 547 or port 546)"; static const char dhcpv6_outbound_filter[] = "outbound and ip6 and udp and (port 547 or port 546)"; -/** Tags for different events, so we can triiger only one type */ +/** Tags for different events, so we can trigger only one type */ static const char packet_handler_tag[] = "PacketHandler"; static const char cache_counter_updater_tag[] = "CacheCounterUpdater"; static const char keepalive_tag[] = "Keepalive"; @@ -570,8 +571,33 @@ int sock_mgr_suspend_packet_handler() syslog(LOG_ALERT, "Packet handlers are already suspended"); return -1; } + std::vector suspended_event_mgrs; for (const auto &entry : sock_map) { - entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); + event_mgr *event_mgr_ptr = entry.second.event_mgr_ptr; + int suspend_result = event_mgr_ptr->suspend_all_events(packet_handler_tag); + if (suspend_result < 0) { + bool rollback_failed = suspend_result < -1; + for (event_mgr *suspended_event_mgr : suspended_event_mgrs) { + if (suspended_event_mgr->resume_all_events(packet_handler_tag) < 0) { + rollback_failed = true; + } + } + if (rollback_failed) { + set_packet_handlers_enabled(false); + for (const auto &rollback_entry : sock_map) { + rollback_entry.second.event_mgr_ptr->suspend_all_events(packet_handler_tag); + } + try { + packet_handler_quiesce_lock = + std::unique_lock(packet_handler_quiesce_mutex); + } catch (const std::system_error &e) { + syslog(LOG_ALERT, "Failed to quiesce packet handlers after suspend rollback failure: %s", + e.what()); + } + } + return -1; + } + suspended_event_mgrs.push_back(event_mgr_ptr); } set_packet_handlers_enabled(false); try { @@ -840,9 +866,9 @@ void sock_mgr_update_db_counters(const socket_counters_t &counters_by_socket) std::string table_name = construct_counter_db_table_key(ifname, info.is_v6); mCountersDbPtr->hset(table_name, info.is_rx ? "RX" : "TX", value); } - syslog_debug(LOG_INFO, "Processing cache counter entry of %sfor downstream vlan %s", + syslog_debug(LOG_INFO, "Processing cache counter entry of %s for downstream vlan %s", all_ifname.c_str(), downstream_ifname.c_str()); - syslog_debug(LOG_INFO, "Skipped aggregated device counter entry of %sfor downstream vlan %s", + syslog_debug(LOG_INFO, "Skipped aggregated device counter entry of %s for downstream vlan %s", all_skipped_ifname.c_str(), downstream_ifname.c_str()); } }