diff --git a/esphome/components/esp32/hal.h b/esphome/components/esp32/hal.h index 2180f07f6c..c01dcc858e 100644 --- a/esphome/components/esp32/hal.h +++ b/esphome/components/esp32/hal.h @@ -9,6 +9,7 @@ #include #include +#include "esphome/core/main_task.h" #include "esphome/core/time_conversion.h" #ifndef PROGMEM @@ -31,6 +32,10 @@ __attribute__((always_inline)) inline bool in_isr_context() { return xPortInIsrC // NOLINTNEXTLINE(readability-redundant-declaration) extern "C" int64_t esp_timer_get_time(void); +/// Before setup() stores the handle every task counts as background; the wake is a no-op until then. +__attribute__((always_inline)) inline bool is_main_loop_thread() { + return xTaskGetCurrentTaskHandle() == esphome_main_task_handle; +} __attribute__((always_inline)) inline void yield() { vPortYield(); } __attribute__((always_inline)) inline void delay(uint32_t ms) { vTaskDelay(ms / portTICK_PERIOD_MS); } __attribute__((always_inline)) inline uint32_t micros() { return static_cast(esp_timer_get_time()); } diff --git a/esphome/components/esp8266/hal.h b/esphome/components/esp8266/hal.h index f3b33da692..1a54065541 100644 --- a/esphome/components/esp8266/hal.h +++ b/esphome/components/esp8266/hal.h @@ -4,6 +4,7 @@ #include #include +#include #include #include @@ -40,6 +41,9 @@ void delay_microseconds_safe(uint32_t us); /// which is ISR-safe) so this helper is unused on this platform. __attribute__((always_inline)) inline bool in_isr_context() { return false; } +/// SDK SYS callbacks and interrupt handlers cannot yield, so they are not the loop context. +__attribute__((always_inline)) inline bool is_main_loop_thread() { return can_yield(); } + __attribute__((always_inline)) inline void yield() { ::yield(); } __attribute__((always_inline)) inline uint32_t micros() { return static_cast(::micros()); } void delay(uint32_t ms); diff --git a/esphome/components/host/hal.h b/esphome/components/host/hal.h index 12abf6684d..c767632772 100644 --- a/esphome/components/host/hal.h +++ b/esphome/components/host/hal.h @@ -3,6 +3,7 @@ #ifdef USE_HOST #include +#include #include #define IRAM_ATTR @@ -16,6 +17,16 @@ namespace esphome { /// Host has no ISR concept. __attribute__((always_inline)) inline bool in_isr_context() { return false; } +inline pthread_t &main_loop_thread() { + static pthread_t thread{}; + return thread; +} + +/// arch_init() stores the thread before any component runs, so no unset case exists. +__attribute__((always_inline)) inline bool is_main_loop_thread() { + return pthread_equal(pthread_self(), main_loop_thread()) != 0; +} + __attribute__((always_inline)) inline void yield() { ::sched_yield(); } void delay(uint32_t ms); @@ -25,7 +36,7 @@ uint64_t millis_64(); void delayMicroseconds(uint32_t us); // NOLINT(readability-identifier-naming) uint32_t arch_get_cpu_cycle_count(); -__attribute__((always_inline)) inline void arch_init() {} +__attribute__((always_inline)) inline void arch_init() { main_loop_thread() = pthread_self(); } __attribute__((always_inline)) inline void arch_feed_wdt() {} __attribute__((always_inline)) inline uint32_t arch_get_cpu_freq_hz() { return 1000000000U; } diff --git a/esphome/components/libretiny/hal.h b/esphome/components/libretiny/hal.h index 48b94a5214..838536b866 100644 --- a/esphome/components/libretiny/hal.h +++ b/esphome/components/libretiny/hal.h @@ -8,6 +8,7 @@ #include #include +#include "esphome/core/main_task.h" #include "esphome/core/time_64.h" // IRAM_ATTR places a function in executable RAM so it is callable from an @@ -84,6 +85,10 @@ __attribute__((always_inline)) inline bool in_isr_context() { #endif } +/// Before setup() stores the handle every task counts as background; the wake is a no-op until then. +__attribute__((always_inline)) inline bool is_main_loop_thread() { + return xTaskGetCurrentTaskHandle() == esphome_main_task_handle; +} __attribute__((always_inline)) inline void yield() { ::yield(); } __attribute__((always_inline)) inline void delay(uint32_t ms) { ::delay(ms); } __attribute__((always_inline)) inline uint32_t micros() { return static_cast(::micros()); } diff --git a/esphome/components/rp2/hal.h b/esphome/components/rp2/hal.h index ec46937bab..2a538f1571 100644 --- a/esphome/components/rp2/hal.h +++ b/esphome/components/rp2/hal.h @@ -4,6 +4,8 @@ #include +#include + #include "esphome/core/time_conversion.h" #define IRAM_ATTR __attribute__((noinline, long_call, section(".time_critical"))) @@ -40,6 +42,9 @@ __attribute__((always_inline)) inline bool in_isr_context() { return ipsr != 0; } +/// arduino-pico runs setup() and loop() on core 0; ESPHome never uses core 1. +__attribute__((always_inline)) inline bool is_main_loop_thread() { return !in_isr_context() && get_core_num() == 0; } + __attribute__((always_inline)) inline void yield() { ::yield(); } __attribute__((always_inline)) inline void delay(uint32_t ms) { ::delay(ms); } __attribute__((always_inline)) inline uint32_t micros() { return static_cast(::micros()); } diff --git a/esphome/components/zephyr/hal.cpp b/esphome/components/zephyr/hal.cpp index 10e8340a40..96c7df896c 100644 --- a/esphome/components/zephyr/hal.cpp +++ b/esphome/components/zephyr/hal.cpp @@ -23,6 +23,7 @@ static const device *const WDT = DEVICE_DT_GET(DT_ALIAS(watchdog0)); // components/zephyr/hal.h. void arch_init() { + main_loop_thread() = k_current_get(); #ifdef CONFIG_WATCHDOG if (device_is_ready(WDT)) { static wdt_timeout_cfg wdt_config{}; diff --git a/esphome/components/zephyr/hal.h b/esphome/components/zephyr/hal.h index 11994b68b7..8384225d33 100644 --- a/esphome/components/zephyr/hal.h +++ b/esphome/components/zephyr/hal.h @@ -17,6 +17,14 @@ namespace esphome { /// Zephyr/nRF52: not currently consulted — wake path is platform-specific. __attribute__((always_inline)) inline bool in_isr_context() { return false; } +inline k_tid_t &main_loop_thread() { + static k_tid_t thread = nullptr; + return thread; +} + +/// arch_init() stores the thread before any component runs, so no unset case exists. +__attribute__((always_inline)) inline bool is_main_loop_thread() { return k_current_get() == main_loop_thread(); } + __attribute__((always_inline)) inline void yield() { ::k_yield(); } __attribute__((always_inline)) inline void delay(uint32_t ms) { ::k_msleep(ms); } __attribute__((always_inline)) inline uint32_t micros() { return k_ticks_to_us_floor32(k_uptime_ticks()); } diff --git a/esphome/components/zigbee/time/zigbee_time_esp32.cpp b/esphome/components/zigbee/time/zigbee_time_esp32.cpp index 1bc0923a1f..909c714c5e 100644 --- a/esphome/components/zigbee/time/zigbee_time_esp32.cpp +++ b/esphome/components/zigbee/time/zigbee_time_esp32.cpp @@ -1,7 +1,6 @@ #include "zigbee_time_esp32.h" #if defined(USE_ZIGBEE) && defined(USE_ESP32) && defined(USE_TIME) #include "esphome/core/log.h" -#include "esphome/core/application.h" namespace esphome::zigbee { @@ -105,7 +104,6 @@ void ZigbeeTime::set_epoch_time(uint32_t utc) { ESP_LOGV(TAG, "Setting device time to UTC: %u", static_cast(utc)); this->synchronize_epoch_(utc); }); - App.wake_loop_threadsafe(); } void ZigbeeTime::dump_config() { diff --git a/esphome/components/zigbee/time/zigbee_time_zephyr.cpp b/esphome/components/zigbee/time/zigbee_time_zephyr.cpp index f0ef3f2af2..7f0a5f979c 100644 --- a/esphome/components/zigbee/time/zigbee_time_zephyr.cpp +++ b/esphome/components/zigbee/time/zigbee_time_zephyr.cpp @@ -1,7 +1,6 @@ #include "zigbee_time_zephyr.h" #if defined(USE_ZIGBEE) && defined(USE_NRF52) && defined(USE_TIME) #include "esphome/core/log.h" -#include "esphome/core/application.h" namespace esphome::zigbee { @@ -48,7 +47,6 @@ void ZigbeeTime::set_epoch_time(uint32_t epoch) { this->synchronize_epoch_(epoch); this->has_time_ = true; }); - App.wake_loop_threadsafe(); } void ZigbeeTime::zcl_device_cb_(zb_bufid_t bufid) { diff --git a/esphome/components/zigbee/zigbee_esp32.cpp b/esphome/components/zigbee/zigbee_esp32.cpp index 8b7b59acb0..89417e748d 100644 --- a/esphome/components/zigbee/zigbee_esp32.cpp +++ b/esphome/components/zigbee/zigbee_esp32.cpp @@ -53,7 +53,6 @@ void ZigbeeComponent::esp_zigbee_alarm_bdb_commissioning(ezb_bdb_comm_mode_mask_ if (!esp_zigbee_lock_acquire(10 / portTICK_PERIOD_MS)) { global_zigbee->set_timeout(COMMISSIONING_RETRY_TIMEOUT_ID, 100, [mode]() { ZigbeeComponent::esp_zigbee_alarm_bdb_commissioning(mode); }); - App.wake_loop_threadsafe(); return; } if (ezb_bdb_start_top_level_commissioning(mode) != EZB_ERR_NONE) { @@ -92,7 +91,6 @@ bool ZigbeeComponent::app_signal_handler(const ezb_app_signal_t *app_signal) { global_zigbee->set_timeout(COMMISSIONING_RETRY_TIMEOUT_ID, 1000, []() { ZigbeeComponent::esp_zigbee_alarm_bdb_commissioning(EZB_BDB_MODE_INITIALIZATION); }); - App.wake_loop_threadsafe(); } } break; case EZB_BDB_SIGNAL_STEERING: { @@ -118,7 +116,6 @@ bool ZigbeeComponent::app_signal_handler(const ezb_app_signal_t *app_signal) { ZigbeeComponent::esp_zigbee_alarm_bdb_commissioning(EZB_BDB_MODE_NETWORK_STEERING); }); } - App.wake_loop_threadsafe(); } } break; case EZB_ZDO_SIGNAL_LEAVE: { diff --git a/esphome/components/zigbee/zigbee_zephyr.cpp b/esphome/components/zigbee/zigbee_zephyr.cpp index f1eadbe6f0..1b2d7392df 100644 --- a/esphome/components/zigbee/zigbee_zephyr.cpp +++ b/esphome/components/zigbee/zigbee_zephyr.cpp @@ -1,7 +1,6 @@ #include "zigbee_zephyr.h" #if defined(USE_ZIGBEE) && defined(USE_NRF52) #include "esphome/core/log.h" -#include "esphome/core/application.h" #include #include #include "esphome/core/hal.h" @@ -120,8 +119,6 @@ void ZigbeeComponent::zcl_device_cb(zb_bufid_t bufid) { /* Set default response value. */ p_device_cb_param->status = RET_OK; - App.wake_loop_threadsafe(); - // endpoints are enumerated from 1 if (global_zigbee->callbacks_.size() >= endpoint) { const auto &cb = global_zigbee->callbacks_[endpoint - 1]; @@ -138,7 +135,6 @@ void ZigbeeComponent::on_join_(bool factory_new) { ESP_LOGD(TAG, "Joined the network"); this->join_cb_.call(factory_new); }); - App.wake_loop_threadsafe(); } void ZigbeeComponent::on_start_() { @@ -146,7 +142,6 @@ void ZigbeeComponent::on_start_() { ESP_LOGD(TAG, "Started zigbee stack"); this->start_cb_.call(); }); - App.wake_loop_threadsafe(); } #ifdef USE_ZIGBEE_WIPE_ON_BOOT @@ -199,6 +194,7 @@ void ZigbeeComponent::setup() { zigbee_configure_sleepy_behavior(this->sleepy_); #endif zigbee_enable(); + this->disable_loop(); } #ifdef ESPHOME_LOG_HAS_CONFIG @@ -263,7 +259,10 @@ static void send_attribute_report(zb_bufid_t bufid, zb_uint16_t cmd_id) { zb_buf_free(bufid); } -void ZigbeeComponent::force_report() { this->force_report_ = true; } +void ZigbeeComponent::force_report() { + this->force_report_ = true; + this->enable_loop_soon_any_context(); +} void ZigbeeComponent::add_radio_sleep_time_ms(uint32_t ms) { this->radio_sleep_remainder_ += ms; @@ -277,6 +276,7 @@ void ZigbeeComponent::loop() { this->force_report_ = false; zb_buf_get_out_delayed_ext(send_attribute_report, 0, 0); } + this->disable_loop(); } void ZigbeeComponent::factory_reset() { diff --git a/esphome/core/scheduler.cpp b/esphome/core/scheduler.cpp index 5acd2eec27..5d0b96c38f 100644 --- a/esphome/core/scheduler.cpp +++ b/esphome/core/scheduler.cpp @@ -138,79 +138,84 @@ void HOT Scheduler::set_timer_common_(Component *component, SchedulerItem::Type } // Take lock early to protect scheduler_item_pool_head_ access - LockGuard guard{this->lock_}; + { + LockGuard guard{this->lock_}; - // Create and populate the scheduler item - SchedulerItem *item = this->get_item_from_pool_locked_(); - // SELF_POINTER items store the source name (owning script) in the union slot instead of a component. - if (name_type == NameType::SELF_POINTER) { - item->source_name = source; - } else { - item->component = component; - } - item->set_name(name_type, static_name, hash_or_id); - item->type = type; - // Use destroy + placement-new instead of move-assignment. - // GCC's std::function::operator=(function&&) does a full swap dance even when the - // target is empty. Since recycled/new items always have an empty callback, we can - // destroy the empty one (no-op) and move-construct directly, saving ~40 bytes of - // swap/destructor code on Xtensa. - item->callback.~function(); - new (&item->callback) std::function(std::move(func)); - // Reset remove flag - recycled items may have been cancelled (remove=true) in previous use - this->set_item_removed_(item, false); + // Create and populate the scheduler item + SchedulerItem *item = this->get_item_from_pool_locked_(); + // SELF_POINTER items store the source name (owning script) in the union slot instead of a component. + if (name_type == NameType::SELF_POINTER) { + item->source_name = source; + } else { + item->component = component; + } + item->set_name(name_type, static_name, hash_or_id); + item->type = type; + // Use destroy + placement-new instead of move-assignment. + // GCC's std::function::operator=(function&&) does a full swap dance even when the + // target is empty. Since recycled/new items always have an empty callback, we can + // destroy the empty one (no-op) and move-construct directly, saving ~40 bytes of + // swap/destructor code on Xtensa. + item->callback.~function(); + new (&item->callback) std::function(std::move(func)); + // Reset remove flag - recycled items may have been cancelled (remove=true) in previous use + this->set_item_removed_(item, false); - // Determine target container: defer_queue_ for deferred items, to_add_ for everything else. - // Using a pointer lets both paths share the cancel + push_back epilogue. - auto *target = &this->to_add_; + // Determine target container: defer_queue_ for deferred items, to_add_ for everything else. + // Using a pointer lets both paths share the cancel + push_back epilogue. + auto *target = &this->to_add_; #ifndef ESPHOME_THREAD_SINGLE - // Special handling for defer() (delay = 0, type = TIMEOUT) - // Single-core platforms don't need thread-safe defer handling - if (delay == 0 && type == SchedulerItem::TIMEOUT) { - // Put in defer queue for guaranteed FIFO execution - target = &this->defer_queue_; - } else + // Special handling for defer() (delay = 0, type = TIMEOUT) + // Single-core platforms don't need thread-safe defer handling + if (delay == 0 && type == SchedulerItem::TIMEOUT) { + // Put in defer queue for guaranteed FIFO execution + target = &this->defer_queue_; + } else #endif /* not ESPHOME_THREAD_SINGLE */ - { - // Only non-defer items need a timestamp for scheduling - const uint64_t now_64 = millis_64(); + { + // Only non-defer items need a timestamp for scheduling + const uint64_t now_64 = millis_64(); - // Type-specific setup - if (type == SchedulerItem::INTERVAL) { - item->interval = delay; - // first execution happens immediately after a random smallish offset - uint32_t offset = this->calculate_interval_offset_(delay); - item->set_next_execution(now_64 + offset); + // Type-specific setup + if (type == SchedulerItem::INTERVAL) { + item->interval = delay; + // first execution happens immediately after a random smallish offset + uint32_t offset = this->calculate_interval_offset_(delay); + item->set_next_execution(now_64 + offset); #ifdef ESPHOME_LOG_HAS_VERBOSE - SchedulerNameLog name_log; - ESP_LOGV(TAG, "Scheduler interval for %s is %" PRIu32 "ms, offset %" PRIu32 "ms", - name_log.format(name_type, static_name, hash_or_id), delay, offset); + SchedulerNameLog name_log; + ESP_LOGV(TAG, "Scheduler interval for %s is %" PRIu32 "ms, offset %" PRIu32 "ms", + name_log.format(name_type, static_name, hash_or_id), delay, offset); #endif - } else { - item->interval = 0; - item->set_next_execution(now_64 + delay); - } + } else { + item->interval = 0; + item->set_next_execution(now_64 + delay); + } #ifdef ESPHOME_DEBUG_SCHEDULER - this->debug_log_timer_(item, name_type, static_name, hash_or_id, delay, now_64); + this->debug_log_timer_(item, name_type, static_name, hash_or_id, delay, now_64); #endif /* ESPHOME_DEBUG_SCHEDULER */ - } + } - // Common epilogue: atomic cancel-and-add (unless skip_cancel is true or anonymous) - // Anonymous items (STATIC_STRING with nullptr) can never match anything, so skip the scan. - if (!skip_cancel && (name_type != NameType::STATIC_STRING || static_name != nullptr)) { - this->cancel_item_locked_(component, name_type, static_name, hash_or_id, type, /* find_first= */ true); - } - target->push_back(item); - if (target == &this->to_add_) { - this->to_add_count_increment_locked_(); - } + // Common epilogue: atomic cancel-and-add (unless skip_cancel is true or anonymous) + // Anonymous items (STATIC_STRING with nullptr) can never match anything, so skip the scan. + if (!skip_cancel && (name_type != NameType::STATIC_STRING || static_name != nullptr)) { + this->cancel_item_locked_(component, name_type, static_name, hash_or_id, type, /* find_first= */ true); + } + target->push_back(item); + if (target == &this->to_add_) { + this->to_add_count_increment_locked_(); + } #ifndef ESPHOME_THREAD_SINGLE - else { - this->defer_count_increment_locked_(); - } + else { + this->defer_count_increment_locked_(); + } #endif + } + // A background insertion may shorten a sleep whose deadline was already computed. + if (!is_main_loop_thread()) [[unlikely]] + wake_scheduler_threadsafe(); } void HOT Scheduler::set_timeout(Component *component, const char *name, uint32_t timeout, diff --git a/esphome/core/wake/wake_esp8266.h b/esphome/core/wake/wake_esp8266.h index 73b7a38a35..2f7ffc292f 100644 --- a/esphome/core/wake/wake_esp8266.h +++ b/esphome/core/wake/wake_esp8266.h @@ -11,10 +11,7 @@ namespace esphome { /// Inline implementation — IRAM callers inline this directly. -inline void ESPHOME_ALWAYS_INLINE wake_loop_impl() { - // Set the wake-requested flag BEFORE esp_schedule so the consumer is - // guaranteed to see it on its next gate check. - wake_request_set(); +inline void ESPHOME_ALWAYS_INLINE wake_scheduler_impl() { // Skip the post when a wake was already signalled and not yet consumed by // wakeable_delay(): esp_schedule() -> ets_post() can enter SDK WiFi pm code, // which must not be poked per-byte from the software serial RX ISR (see @@ -26,11 +23,19 @@ inline void ESPHOME_ALWAYS_INLINE wake_loop_impl() { esp_schedule(); } +inline void ESPHOME_ALWAYS_INLINE wake_loop_impl() { + // Set the wake-requested flag BEFORE esp_schedule so the consumer is + // guaranteed to see it on its next gate check. + wake_request_set(); + wake_scheduler_impl(); +} + /// IRAM_ATTR entry point for ISR callers — defined in wake_esp8266.cpp. void wake_loop_any_context(); /// Non-ISR: always inline. inline void wake_loop_threadsafe() { wake_loop_impl(); } +inline void wake_scheduler_threadsafe() { wake_scheduler_impl(); } /// ISR-safe: no task_woken arg because ESP8266 has no FreeRTOS. Caller must be IRAM_ATTR. inline void ESPHOME_ALWAYS_INLINE wake_loop_isrsafe() { wake_loop_impl(); } diff --git a/esphome/core/wake/wake_freertos.cpp b/esphome/core/wake/wake_freertos.cpp index 458ef51f89..02c628718e 100644 --- a/esphome/core/wake/wake_freertos.cpp +++ b/esphome/core/wake/wake_freertos.cpp @@ -30,9 +30,6 @@ void IRAM_ATTR wake_loop_any_context() { wake_main_task_any_context(); } } // namespace esphome -extern "C" void esphome_wake_loop_threadsafe() { - esphome::wake_request_set(); - esphome_main_task_notify(); -} +extern "C" void esphome_wake_loop_threadsafe() { esphome::wake_loop_threadsafe(); } #endif // USE_ESP32 || USE_LIBRETINY diff --git a/esphome/core/wake/wake_freertos.h b/esphome/core/wake/wake_freertos.h index 16afa38fda..1df7d38765 100644 --- a/esphome/core/wake/wake_freertos.h +++ b/esphome/core/wake/wake_freertos.h @@ -34,9 +34,11 @@ __attribute__((always_inline)) inline void wake_main_task_any_context() { void wake_loop_isrsafe(BaseType_t *px_higher_priority_task_woken); void wake_loop_any_context(); +inline void wake_scheduler_threadsafe() { esphome_main_task_notify(); } + inline void wake_loop_threadsafe() { wake_request_set(); - esphome_main_task_notify(); + wake_scheduler_threadsafe(); } namespace internal { diff --git a/esphome/core/wake/wake_host.cpp b/esphome/core/wake/wake_host.cpp index 8cb382a77e..9fc1c71703 100644 --- a/esphome/core/wake/wake_host.cpp +++ b/esphome/core/wake/wake_host.cpp @@ -123,13 +123,12 @@ void wakeable_delay(uint32_t ms) { if (ms == 0) [[unlikely]] { yield(); } - // A socket woke select() early — open the component-phase gate so the - // owning component's loop() drains the data on this tick rather than - // waiting up to loop_interval_ ms. Idempotent if wake_loop_threadsafe() - // already set the flag (wake socket fired); required when an application - // socket fired and nothing else set the flag. + // Application sockets need the component phase to drain queued work. + // The internal wake socket may only be signaling new scheduler work. if (ret > 0) { - wake_request_set(); + const bool only_wake_socket = ret == 1 && g_wake_socket_fd >= 0 && FD_ISSET(g_wake_socket_fd, &g_read_fds); + if (!only_wake_socket) + wake_request_set(); } return; } @@ -146,16 +145,20 @@ void wakeable_delay(uint32_t ms) { } } // namespace internal -void wake_loop_threadsafe() { - // Set flag before sending so the consumer's gate check on the next loop() - // entry observes the wake regardless of select() scheduling. - wake_request_set(); +void wake_scheduler_threadsafe() { if (internal::g_wake_socket_fd >= 0) { const char dummy = 1; ::send(internal::g_wake_socket_fd, &dummy, 1, 0); } } +void wake_loop_threadsafe() { + // Set flag before sending so the consumer's gate check on the next loop() + // entry observes the wake regardless of select() scheduling. + wake_request_set(); + wake_scheduler_threadsafe(); +} + void wake_setup() { // Create UDP socket for wake notifications. internal::g_wake_socket_fd = ::socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); diff --git a/esphome/core/wake/wake_host.h b/esphome/core/wake/wake_host.h index 9756ed4c39..8c58797df2 100644 --- a/esphome/core/wake/wake_host.h +++ b/esphome/core/wake/wake_host.h @@ -13,6 +13,7 @@ namespace esphome { /// Host: wakes select() via UDP loopback socket. Defined in wake_host.cpp. void wake_loop_threadsafe(); +void wake_scheduler_threadsafe(); /// Register a socket file descriptor with the host select() loop. Not /// thread-safe — main loop only. Returns false if fd is invalid or diff --git a/esphome/core/wake/wake_rp2.h b/esphome/core/wake/wake_rp2.h index 715e5aca0c..4f5439dbb8 100644 --- a/esphome/core/wake/wake_rp2.h +++ b/esphome/core/wake/wake_rp2.h @@ -11,12 +11,16 @@ namespace esphome { +inline void wake_scheduler_threadsafe() { + g_main_loop_woke = true; + __sev(); +} + inline void wake_loop_any_context() { // Set the wake-requested flag BEFORE the SEV so the consumer is guaranteed // to see it on its next gate check. wake_request_set(); - g_main_loop_woke = true; - __sev(); + wake_scheduler_threadsafe(); } inline void wake_loop_threadsafe() { wake_loop_any_context(); } diff --git a/esphome/core/wake/wake_zephyr.cpp b/esphome/core/wake/wake_zephyr.cpp index 577d53f5d9..2c46fbc951 100644 --- a/esphome/core/wake/wake_zephyr.cpp +++ b/esphome/core/wake/wake_zephyr.cpp @@ -21,9 +21,11 @@ K_SEM_DEFINE(esphome_wake_sem, 0, 1); // NOLINTNEXTLINE(cppcoreguidelines-avoid-non-const-global-variables) volatile uint8_t g_wake_requested = 0; +void wake_scheduler_threadsafe() { k_sem_give(&esphome_wake_sem); } + void wake_loop_threadsafe() { wake_request_set(); - k_sem_give(&esphome_wake_sem); + wake_scheduler_threadsafe(); } namespace internal { diff --git a/esphome/core/wake/wake_zephyr.h b/esphome/core/wake/wake_zephyr.h index c89cfc68e9..04cbfb230b 100644 --- a/esphome/core/wake/wake_zephyr.h +++ b/esphome/core/wake/wake_zephyr.h @@ -11,6 +11,7 @@ namespace esphome { /// Zephyr: wakes the main loop via k_sem_give(). Thread- and ISR-safe. /// Defined in wake_zephyr.cpp. void wake_loop_threadsafe(); +void wake_scheduler_threadsafe(); inline void wake_loop_any_context() { wake_loop_threadsafe(); } diff --git a/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.cpp b/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.cpp index b58f1c9adc..71182a46ad 100644 --- a/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.cpp +++ b/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.cpp @@ -16,4 +16,16 @@ void WakeTestComponent::start_async_wake() { }).detach(); } +void WakeTestComponent::start_async_timeout(uint32_t delay_ms) { + const uint32_t start_time = millis(); + std::thread([this, start_time, delay_ms] { + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + const int start_loop_count = this->get_loop_count(); + this->set_timeout(delay_ms, [this, start_time, start_loop_count, delay_ms] { + ESP_LOGI(TAG, "SCHEDULER_WAKE_RESULT delay=%u elapsed=%u loop_delta=%d", static_cast(delay_ms), + static_cast(millis() - start_time), this->get_loop_count() - start_loop_count); + }); + }).detach(); +} + } // namespace esphome::wake_test_component diff --git a/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.h b/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.h index c8e4e0a89f..a1f06f8998 100644 --- a/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.h +++ b/tests/integration/fixtures/external_components/wake_test_component/wake_test_component.h @@ -18,6 +18,10 @@ class WakeTestComponent : public Component { // loop_interval_ has been raised high enough to gate it off otherwise. void start_async_wake(); + // Spawn a detached thread that inserts a timeout (or a defer when delay_ms + // is 0) while the main loop sleeps, and report how soon it ran. + void start_async_timeout(uint32_t delay_ms); + float get_setup_priority() const override { return setup_priority::DATA; } protected: diff --git a/tests/integration/fixtures/scheduler_background_wake.yaml b/tests/integration/fixtures/scheduler_background_wake.yaml new file mode 100644 index 0000000000..78e6864eb2 --- /dev/null +++ b/tests/integration/fixtures/scheduler_background_wake.yaml @@ -0,0 +1,31 @@ +esphome: + name: scheduler-background-wake + on_boot: + priority: -100 + then: + - lambda: |- + App.set_loop_interval(5000); + +host: +api: +logger: + level: INFO + +external_components: + - source: + type: local + path: EXTERNAL_COMPONENT_PATH + components: [wake_test_component] + +wake_test_component: + id: wake_test + +button: + - platform: template + name: Start Scheduler Timeout + on_press: + - lambda: id(wake_test)->start_async_timeout(100); + - platform: template + name: Start Scheduler Defer + on_press: + - lambda: id(wake_test)->start_async_timeout(0); diff --git a/tests/integration/test_scheduler_background_wake.py b/tests/integration/test_scheduler_background_wake.py new file mode 100644 index 0000000000..fa7d34f494 --- /dev/null +++ b/tests/integration/test_scheduler_background_wake.py @@ -0,0 +1,84 @@ +"""A background scheduler insertion must interrupt the loop's current sleep. + +Covers both set_timeout and defer: a zero-delay insert from another thread +takes the separate defer queue on multi-threaded builds. +""" + +from __future__ import annotations + +import asyncio +from pathlib import Path +import re + +from aioesphomeapi import ButtonInfo +import pytest + +from .state_utils import require_entity +from .types import APIClientConnectedFactory, RunCompiledFunction + + +@pytest.mark.asyncio +async def test_scheduler_background_wake( + yaml_config: str, + run_compiled: RunCompiledFunction, + api_client_connected: APIClientConnectedFactory, +) -> None: + external_components_path = str( + Path(__file__).parent / "fixtures" / "external_components" + ) + yaml_config = yaml_config.replace( + "EXTERNAL_COMPONENT_PATH", external_components_path + ) + + loop = asyncio.get_running_loop() + results: dict[int, asyncio.Future[tuple[int, int]]] = { + 100: loop.create_future(), + 0: loop.create_future(), + } + + def on_log_line(line: str) -> None: + match = re.search( + r"SCHEDULER_WAKE_RESULT delay=(\d+) elapsed=(\d+) loop_delta=(-?\d+)", + line, + ) + if match is None: + return + result = results[int(match.group(1))] + if not result.done(): + result.set_result((int(match.group(2)), int(match.group(3)))) + + async with ( + run_compiled(yaml_config, line_callback=on_log_line), + api_client_connected() as client, + ): + device_info = await client.device_info() + assert device_info is not None + assert device_info.name == "scheduler-background-wake" + entities, _ = await client.list_entities_services() + for object_id, delay in ( + ("start_scheduler_timeout", 100), + ("start_scheduler_defer", 0), + ): + start_button = require_entity( + entities, + object_id, + ButtonInfo, + description=f"{object_id} button", + ) + client.button_command(start_button.key) + + try: + elapsed, loop_delta = await asyncio.wait_for( + results[delay], timeout=10.0 + ) + except TimeoutError: + pytest.fail(f"background insert with delay={delay} did not fire") + + assert elapsed < 1000, ( + f"background insert with delay={delay} should interrupt the " + f"five-second sleep; it fired after {elapsed}ms" + ) + assert loop_delta == 0, ( + f"scheduler-only wake must not run component loops; observed " + f"loop_delta={loop_delta} for delay={delay}" + )