From f4a055f342435dfbba5913efad74bf408e36f1c8 Mon Sep 17 00:00:00 2001 From: "J. Nick Koston" Date: Fri, 24 Apr 2026 02:55:23 -0500 Subject: [PATCH] [web_server_idf] Defer SSE session adoption/priming to the main loop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit AsyncEventSource::handleRequest() runs on the ESP-IDF httpd task. Prior to this change it mutated sessions_, event_buffer_, deferred_queue_, and entities_iterator_ directly: appending the new response to sessions_, sending the initial ping/config/sorting_groups through try_send_nodefer() (which writes event_buffer_), and calling entities_iterator_.begin(). The main loop reads/writes the same fields from WebServer::loop() on the main task, so every new SSE connect raced with any in-flight loop tick. The httpd path now only performs the HTTP-level setup that requires the live httpd_req_t (headers, initial chunk, sess_ctx/free_ctx, fd_, send override) and parks the response on pending_sessions_ under a mutex. The main loop checks a std::atomic fast-path gate per tick; when set it swaps the pending list out under the lock and then — outside the lock — pushes into sessions_, invokes on_connect_, and calls prime_() which performs the initial sends and starts entities_iterator_. sessions_, event_buffer_, deferred_queue_, and entities_iterator_ are now mutated exclusively from the main loop. fd_ remains the only cross-thread signal (std::atomic cleared by destroy() on the httpd/tcpip task). --- .../web_server_idf/web_server_idf.cpp | 55 +++++++++++++++---- .../web_server_idf/web_server_idf.h | 11 ++++ 2 files changed, 56 insertions(+), 10 deletions(-) diff --git a/esphome/components/web_server_idf/web_server_idf.cpp b/esphome/components/web_server_idf/web_server_idf.cpp index 8f464ae912..3610f4cdfa 100644 --- a/esphome/components/web_server_idf/web_server_idf.cpp +++ b/esphome/components/web_server_idf/web_server_idf.cpp @@ -475,21 +475,54 @@ AsyncEventSource::~AsyncEventSource() { for (auto *ses : this->sessions_) { delete ses; // NOLINT(cppcoreguidelines-owning-memory) } + LockGuard guard{this->pending_mutex_}; + for (auto *ses : this->pending_sessions_) { + delete ses; // NOLINT(cppcoreguidelines-owning-memory) + } } void AsyncEventSource::handleRequest(AsyncWebServerRequest *request) { + // Runs on the httpd task. Do only the HTTP-level setup that needs the live httpd_req_t, + // then hand the session off to the main loop for adoption and priming. This keeps + // sessions_, event_buffer_, and deferred_queue_ mutated exclusively from the main loop. // NOLINTNEXTLINE(cppcoreguidelines-owning-memory,clang-analyzer-cplusplus.NewDeleteLeaks) auto *rsp = new AsyncEventSourceResponse(request, this, this->web_server_); - if (this->on_connect_) { - this->on_connect_(rsp); + { + LockGuard guard{this->pending_mutex_}; + this->pending_sessions_.push_back(rsp); + // Release-store so the main loop's acquire-load sees the push_back above. + this->has_pending_sessions_.store(true, std::memory_order_release); } - this->sessions_.push_back(rsp); - // Wake up WebServer::loop() to drain deferred event queues for this client. + // Wake up WebServer::loop() to adopt and prime this client. // Safe from httpd task context via the pending_enable_loop_ flag. this->web_server_->enable_loop_soon_any_context(); } bool AsyncEventSource::loop() { + // Adopt sessions handed off from the httpd task. Fast path: one atomic load per tick + // when nothing is pending. Only take the lock / touch the vector on a real connect. + // Swap under the lock and do the heavy work (on_connect_ callback, initial sends, + // entity iterator start) outside it so they cannot race with httpd handlers. + if (this->has_pending_sessions_.load(std::memory_order_acquire)) { + std::vector incoming; + { + LockGuard guard{this->pending_mutex_}; + incoming.swap(this->pending_sessions_); + this->has_pending_sessions_.store(false, std::memory_order_relaxed); + } + for (auto *rsp : incoming) { + this->sessions_.push_back(rsp); + if (this->on_connect_) { + this->on_connect_(rsp); + } + // Skip priming if the client disconnected before we got here; the cleanup pass below + // will delete it. + if (rsp->fd_.load() != 0) { + rsp->prime_(); + } + } + } + // Clean up dead sessions safely // This follows the ESP-IDF pattern where free_ctx marks resources as dead // and the main loop handles the actual cleanup to avoid race conditions @@ -534,6 +567,9 @@ AsyncEventSourceResponse::AsyncEventSourceResponse(const AsyncWebServerRequest * esphome::web_server_idf::AsyncEventSource *server, esphome::web_server::WebServer *ws) : server_(server), web_server_(ws), entities_iterator_(ws, server) { + // Runs on the httpd task. Only touch state tied to the live httpd_req_t here; the + // main loop will call prime_() later to do the initial sends and start the iterator. + // Writing to event_buffer_ / deferred_queue_ from this task would race with the main loop. httpd_req_t *req = *request; httpd_resp_set_status(req, HTTPD_200); @@ -555,6 +591,11 @@ AsyncEventSourceResponse::AsyncEventSourceResponse(const AsyncWebServerRequest * // Use non-blocking send to prevent watchdog timeouts when TCP buffers are full httpd_sess_set_send_override(this->hd_, this->fd_.load(), nonblocking_send); +} + +void AsyncEventSourceResponse::prime_() { + // Runs on the main loop after AsyncEventSource::loop() adopts this session. + auto *ws = this->web_server_; // Configure reconnect timeout and send config // this should always go through since the tcp send buffer is empty on connect @@ -578,12 +619,6 @@ AsyncEventSourceResponse::AsyncEventSourceResponse(const AsyncWebServerRequest * #endif this->entities_iterator_.begin(ws->include_internal_); - - // just dump them all up-front and take advantage of the deferred queue - // on second thought that takes too long, but leaving the commented code here for debug purposes - // while(!this->entities_iterator_.completed()) { - // this->entities_iterator_.advance(); - //} } void AsyncEventSourceResponse::destroy(void *ptr) { diff --git a/esphome/components/web_server_idf/web_server_idf.h b/esphome/components/web_server_idf/web_server_idf.h index f2931fb507..04a856ae04 100644 --- a/esphome/components/web_server_idf/web_server_idf.h +++ b/esphome/components/web_server_idf/web_server_idf.h @@ -299,6 +299,10 @@ class AsyncEventSourceResponse { AsyncEventSourceResponse(const AsyncWebServerRequest *request, esphome::web_server_idf::AsyncEventSource *server, esphome::web_server::WebServer *ws); + // Sends the initial ping/config/sorting_groups and starts the entity iterator. + // Must be called from the main loop (writes event_buffer_ and touches entities_iterator_). + void prime_(); + void deq_push_back_with_dedup_(void *source, message_generator_t *message_generator); void process_deferred_queue_(); void process_buffer_(); @@ -351,7 +355,14 @@ class AsyncEventSource : public AsyncWebHandler { // Use vector instead of set: SSE sessions are typically 1-5 connections (browsers, dashboards). // Linear search is faster than red-black tree overhead for this small dataset. // Only operations needed: add session, remove session, iterate sessions - no need for sorted order. + // Mutated only from the main loop. std::vector sessions_; + // Sessions constructed on the httpd task wait here until the main loop adopts them. + // All mutations are guarded by pending_mutex_. has_pending_sessions_ is a fast-path + // gate so the per-tick cost in loop() when no connect is pending is one atomic load. + std::vector pending_sessions_; + Mutex pending_mutex_; + std::atomic has_pending_sessions_{false}; connect_handler_t on_connect_{}; esphome::web_server::WebServer *web_server_; };