[core][zigbee] Wake main loop on defer/set_timeout from background task (#20356)

Co-authored-by: J. Nick Koston <nick@home-assistant.io>
This commit is contained in:
luar123
2026-10-09 18:58:33 +00:00
committed by GitHub
co-authored by J. Nick Koston
parent dd8315059c
commit dbad1d4629
24 changed files with 278 additions and 95 deletions
+5
View File
@@ -9,6 +9,7 @@
#include <freertos/FreeRTOS.h>
#include <freertos/task.h>
#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<uint32_t>(esp_timer_get_time()); }
+4
View File
@@ -4,6 +4,7 @@
#include <c_types.h>
#include <core_esp8266_features.h>
#include <coredecls.h>
#include <cstdint>
#include <pgmspace.h>
@@ -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<uint32_t>(::micros()); }
void delay(uint32_t ms);
+12 -1
View File
@@ -3,6 +3,7 @@
#ifdef USE_HOST
#include <cstdint>
#include <pthread.h>
#include <sched.h>
#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; }
+5
View File
@@ -8,6 +8,7 @@
#include <FreeRTOS.h>
#include <task.h>
#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<uint32_t>(::micros()); }
+5
View File
@@ -4,6 +4,8 @@
#include <cstdint>
#include <pico/platform.h>
#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<uint32_t>(::micros()); }
+1
View File
@@ -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{};
+8
View File
@@ -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()); }
@@ -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<unsigned>(utc));
this->synchronize_epoch_(utc);
});
App.wake_loop_threadsafe();
}
void ZigbeeTime::dump_config() {
@@ -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) {
@@ -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: {
+6 -6
View File
@@ -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 <zephyr/settings/settings.h>
#include <zephyr/storage/flash_map.h>
#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() {
+64 -59
View File
@@ -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<void()>(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<void()>(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,
+9 -4
View File
@@ -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(); }
+1 -4
View File
@@ -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
+3 -1
View File
@@ -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 {
+13 -10
View File
@@ -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);
+1
View File
@@ -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
+6 -2
View File
@@ -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(); }
+3 -1
View File
@@ -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 {
+1
View File
@@ -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(); }
@@ -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<unsigned int>(delay_ms),
static_cast<unsigned int>(millis() - start_time), this->get_loop_count() - start_loop_count);
});
}).detach();
}
} // namespace esphome::wake_test_component
@@ -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:
@@ -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);
@@ -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}"
)