mirror of
https://github.com/esphome/esphome.git
synced 2026-09-08 05:56:02 +00:00
Compare commits
22
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d12305e398 | ||
|
|
d34d3994e1 | ||
|
|
9a877a067c | ||
|
|
9ba4477ada | ||
|
|
8434dc5474 | ||
|
|
e0e85db822 | ||
|
|
5d2ddc658c | ||
|
|
c1aa41f276 | ||
|
|
96b1a03ea4 | ||
|
|
95ab3fb4f2 | ||
|
|
18220e0b39 | ||
|
|
011497d6ee | ||
|
|
7089dae3b6 | ||
|
|
745eb30109 | ||
|
|
e36445fa5f | ||
|
|
cb0c2bdaca | ||
|
|
e47247486b | ||
|
|
657116a213 | ||
|
|
3321566cc0 | ||
|
|
d58b37faa1 | ||
|
|
8966567be0 | ||
|
|
20c7dcb1dd |
@@ -37,6 +37,25 @@ CONF_HANDSHAKE_PIN = "handshake_pin"
|
||||
CONF_SDIO_FREQUENCY = "sdio_frequency"
|
||||
CONF_SPI_MODE = "spi_mode"
|
||||
|
||||
# ESP-NOW-over-hosted shim (esp_now_hosted.cpp). esp-hosted proxies esp_wifi.h
|
||||
# but not esp_now.h (espressif/esp-hosted-mcu#19), and esp_wifi_remote injects
|
||||
# the esp_now.h header on the ESP32-P4 host with no implementation, leaving the
|
||||
# esp_now_* symbols undefined at link. On a P4 host, esp_now_hosted.cpp DEFINES
|
||||
# those symbols and forwards each call to the co-processor over esp-hosted's
|
||||
# CustomRpc "peer data transfer" channel, so ESPHome's `espnow` component links
|
||||
# and runs unchanged (proven on a Tab5, 2026-07-20). The .cpp is guarded to
|
||||
# CONFIG_IDF_TARGET_ESP32P4 so it compiles to nothing on hosts with a native
|
||||
# ESP-NOW stack. CustomRpc needs these two host-side Kconfig options. Host
|
||||
# registers 3 handlers (RESP, RECV, SEND); the coprocessor registers 1 (REQ);
|
||||
# we ask for 8 to leave room for other CustomRpc extensions alongside.
|
||||
#
|
||||
# The coprocessor must run the matching custom firmware (a parallel effort in
|
||||
# esphome/esp-hosted-firmware). esp_now_hosted_rpc.h here is the canonical copy
|
||||
# of the wire contract and MUST stay byte-identical to the copy that coprocessor
|
||||
# firmware uses — the packed structs are the on-wire layout, so any divergence
|
||||
# silently corrupts every ESP-NOW frame.
|
||||
_MAX_CUSTOM_MSG_HANDLERS = 8
|
||||
|
||||
# Shared fields for both transport modes
|
||||
BASE_SCHEMA = cv.Schema(
|
||||
{
|
||||
@@ -262,6 +281,23 @@ async def to_code(config: ConfigType) -> None:
|
||||
else:
|
||||
_configure_spi(config)
|
||||
|
||||
# ESP-NOW-over-hosted shim: only the radio-less ESP32-P4 host needs it (see
|
||||
# the note by _MAX_CUSTOM_MSG_HANDLERS). Enabled for every P4 host, not
|
||||
# gated on the `espnow` component being present: the shim is tiny and the
|
||||
# esp_now_* symbols/CustomRpc calls it defines require these Kconfig options
|
||||
# to link whenever esp_now_hosted.cpp compiles (which is on any P4 host), so
|
||||
# coupling the two keeps the build consistent. When `espnow` is absent the
|
||||
# symbols are simply unused and never register a callback at runtime.
|
||||
if esp32.get_esp32_variant() == esp32.VARIANT_ESP32P4:
|
||||
add_define("USE_ESP_NOW_HOSTED")
|
||||
# esp-hosted's CustomRpc ("peer data transfer") path — off by default.
|
||||
esp32.add_idf_sdkconfig_option(
|
||||
"CONFIG_ESP_HOSTED_ENABLE_PEER_DATA_TRANSFER", True
|
||||
)
|
||||
esp32.add_idf_sdkconfig_option(
|
||||
"CONFIG_ESP_HOSTED_MAX_CUSTOM_MSG_HANDLERS", _MAX_CUSTOM_MSG_HANDLERS
|
||||
)
|
||||
|
||||
# Place the transport mempool in PSRAM. Required on memory-tight host
|
||||
# configurations (e.g. P4 with a large LVGL UI) where the internal-RAM
|
||||
# mempool allocation fails at boot with `sdio_mempool_create` assert.
|
||||
|
||||
@@ -0,0 +1,467 @@
|
||||
/*
|
||||
* esp_now_hosted — host-side shim implementing <esp_now.h> over esp-hosted
|
||||
* CustomRpc, so ESPHome's `espnow` component can run on a radio-less host
|
||||
* (e.g. the ESP32-P4) whose radio lives on an esp-hosted co-processor.
|
||||
*
|
||||
* A radio-less host has no native ESP-NOW. esp_wifi_remote INJECTS the full
|
||||
* esp_now.h header (types + declarations) but ships NO implementation, so every
|
||||
* esp_now_* symbol is an undefined reference at link time. This translation
|
||||
* unit provides those definitions; each forwards to the co-processor over
|
||||
* CustomRpc (see esphome/esp-hosted-firmware for the matching coprocessor
|
||||
* handlers). No esp-hosted or esp_wifi_remote source is patched, and there is no
|
||||
* duplicate-symbol clash because nothing else defines these symbols here.
|
||||
*
|
||||
* See esp_now_hosted_rpc.h for the wire protocol.
|
||||
*/
|
||||
|
||||
#include "sdkconfig.h"
|
||||
|
||||
// Only build the shim on the radio-less host. On chips with a native ESP-NOW
|
||||
// stack (S3, C6, …) the real symbols exist and this file must stay empty to
|
||||
// avoid duplicate definitions.
|
||||
#if defined(CONFIG_IDF_TARGET_ESP32P4)
|
||||
|
||||
#include <cstring>
|
||||
|
||||
#include "freertos/FreeRTOS.h"
|
||||
#include "freertos/semphr.h"
|
||||
|
||||
#include "esp_idf_version.h"
|
||||
#include "esp_log.h"
|
||||
#include "esp_timer.h"
|
||||
|
||||
#include <esp_now.h> // injected declarations we are now DEFINING
|
||||
#include <esp_wifi_types.h> // wifi_pkt_rx_ctrl_t, wifi_tx_info_t
|
||||
|
||||
// esp_hosted_misc.h (host) ships WITHOUT an extern "C" guard, so including it
|
||||
// from C++ would give its declarations C++ linkage and the real C symbols in
|
||||
// libesp_hosted would go unresolved at link. Wrap it. (Verified vs
|
||||
// esp_hosted 2.12.9.)
|
||||
extern "C" {
|
||||
#include "esp_hosted_misc.h" // esp_hosted_{send_custom_data,register_custom_callback}
|
||||
}
|
||||
|
||||
#include "esp_now_hosted_rpc.h"
|
||||
|
||||
namespace {
|
||||
|
||||
const char *const TAG = "esp_now_hosted";
|
||||
|
||||
// One outstanding request at a time. ESPHome drives esp_now_* from the main
|
||||
// loop; the matching response and the async RECV/SEND events all arrive on the
|
||||
// single esp-hosted RPC RX thread. Serializing requests keeps the shared
|
||||
// response slot race-free; a sequence number stops a late/stale response from
|
||||
// being mistaken for ours.
|
||||
SemaphoreHandle_t g_req_mutex = nullptr;
|
||||
SemaphoreHandle_t g_resp_sem = nullptr; // given when the matching RESP lands
|
||||
bool g_setup_done = false; // set only after setup fully succeeds
|
||||
uint8_t g_seq = 0;
|
||||
volatile uint8_t g_expect_seq = 0;
|
||||
volatile int32_t g_resp_status = 0;
|
||||
uint8_t g_resp_ret[16];
|
||||
volatile uint16_t g_resp_ret_len = 0;
|
||||
|
||||
// Written from the main loop (register/unregister/deinit), read from the
|
||||
// esp-hosted RX thread (on_recv/on_send). volatile for the same reason the
|
||||
// g_resp_* globals are: force the RX thread to observe an updated pointer
|
||||
// (e.g. a nulling by esp_now_deinit) rather than a cached one.
|
||||
volatile esp_now_recv_cb_t g_recv_cb = nullptr;
|
||||
volatile esp_now_send_cb_t g_send_cb = nullptr;
|
||||
|
||||
// Local mirror of the co-processor's peer table. ESPHome's espnow component
|
||||
// calls esp_now_is_peer_exist() on the main loop for every received frame
|
||||
// (twice) and every send; forwarding each as a blocking RPC round-trip stalls
|
||||
// the loop. The shim is the only path that mutates the co-processor peer table
|
||||
// (add/del/deinit all go through here), so this mirror is authoritative and
|
||||
// esp_now_is_peer_exist() can answer from it with no round-trip.
|
||||
//
|
||||
// esp_now_* are public C symbols: any component or user lambda may call them,
|
||||
// and although ESPHome's espnow touches peers only from the main loop today
|
||||
// (its RX/TX callbacks merely enqueue), the shim cannot rely on that. A short
|
||||
// spinlock keeps the mirror consistent from any task/core, matching native
|
||||
// esp_now_*'s own internal thread-safety. The critical sections are a bounded
|
||||
// (<=20-entry) scan, so they stay tiny. ESP_NOW_MAX_TOTAL_PEER_NUM is 20.
|
||||
constexpr size_t ESP_NOW_HOSTED_MAX_PEERS = 20;
|
||||
uint8_t g_peer_cache[ESP_NOW_HOSTED_MAX_PEERS][6];
|
||||
size_t g_peer_count = 0;
|
||||
portMUX_TYPE g_peer_lock = portMUX_INITIALIZER_UNLOCKED;
|
||||
|
||||
// Caller must hold g_peer_lock.
|
||||
int peer_cache_find_locked(const uint8_t *mac) {
|
||||
for (size_t i = 0; i < g_peer_count; i++) {
|
||||
if (memcmp(g_peer_cache[i], mac, 6) == 0)
|
||||
return static_cast<int>(i);
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
bool peer_cache_contains(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
const bool found = peer_cache_find_locked(mac) >= 0;
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
return found;
|
||||
}
|
||||
|
||||
void peer_cache_add(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
if (peer_cache_find_locked(mac) < 0 && g_peer_count < ESP_NOW_HOSTED_MAX_PEERS)
|
||||
memcpy(g_peer_cache[g_peer_count++], mac, 6);
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
void peer_cache_remove(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
const int idx = peer_cache_find_locked(mac);
|
||||
if (idx >= 0) {
|
||||
g_peer_count--;
|
||||
if (static_cast<size_t>(idx) != g_peer_count) // move the last entry into the gap
|
||||
memcpy(g_peer_cache[idx], g_peer_cache[g_peer_count], 6);
|
||||
}
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
void peer_cache_clear() {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
g_peer_count = 0;
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
// ── CustomRpc event handlers (run on the esp-hosted RPC RX thread) ──────────
|
||||
// Keep them short and non-blocking. In particular they MUST NOT call back into
|
||||
// any esp_now_* shim function: that would try to take g_req_mutex / wait on the
|
||||
// RX thread that delivers the response, and deadlock.
|
||||
|
||||
void on_resp(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
if (len < sizeof(esp_now_hosted_resp_t)) {
|
||||
ESP_LOGW(TAG, "RESP too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *r = reinterpret_cast<const esp_now_hosted_resp_t *>(data);
|
||||
if (r->seq != g_expect_seq) { // late response from a timed-out request (expected)
|
||||
ESP_LOGV(TAG, "dropping stale RESP seq %u (want %u)", r->seq, g_expect_seq);
|
||||
return;
|
||||
}
|
||||
g_resp_status = r->status;
|
||||
uint16_t rl = r->ret_len;
|
||||
if (rl > sizeof(g_resp_ret)) {
|
||||
// Larger than any real opcode return — a likely wire-format drift signal.
|
||||
ESP_LOGW(TAG, "RESP ret_len %u exceeds buffer, clamping (wire drift?)", rl);
|
||||
rl = sizeof(g_resp_ret);
|
||||
}
|
||||
if (len >= sizeof(esp_now_hosted_resp_t) + rl) {
|
||||
memcpy(g_resp_ret, r->ret, rl);
|
||||
} else {
|
||||
// Truncated frame: fail closed. Never hand the caller stale bytes left in
|
||||
// g_resp_ret by a previous response, and don't let request() report a
|
||||
// zeroed payload as success — override the status to an error.
|
||||
ESP_LOGW(TAG, "RESP truncated: claims %u ret bytes, frame too short", rl);
|
||||
rl = 0;
|
||||
g_resp_status = ESP_ERR_INVALID_RESPONSE;
|
||||
}
|
||||
g_resp_ret_len = rl;
|
||||
xSemaphoreGive(g_resp_sem);
|
||||
}
|
||||
|
||||
void on_recv(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
// Read the volatile pointer once: esp_now_unregister_recv_cb()/deinit() (via
|
||||
// the espnow component's disable()) can null it on the main loop between the
|
||||
// guard and the call, which would otherwise turn the call into a null-deref.
|
||||
const esp_now_recv_cb_t cb = g_recv_cb;
|
||||
if (cb == nullptr)
|
||||
return;
|
||||
if (len < sizeof(esp_now_hosted_recv_evt_t)) {
|
||||
ESP_LOGW(TAG, "RECV too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *e = reinterpret_cast<const esp_now_hosted_recv_evt_t *>(data);
|
||||
if (len < sizeof(esp_now_hosted_recv_evt_t) + e->data_len) {
|
||||
ESP_LOGW(TAG, "RECV data_len %u exceeds frame", e->data_len);
|
||||
return;
|
||||
}
|
||||
|
||||
// ESPHome dereferences info->rx_ctrl->{rssi,timestamp}; give it a real one.
|
||||
wifi_pkt_rx_ctrl_t rx_ctrl;
|
||||
memset(&rx_ctrl, 0, sizeof(rx_ctrl));
|
||||
rx_ctrl.rssi = e->rssi;
|
||||
rx_ctrl.channel = e->channel;
|
||||
rx_ctrl.timestamp = static_cast<uint32_t>(esp_timer_get_time());
|
||||
|
||||
esp_now_recv_info_t info;
|
||||
info.src_addr = const_cast<uint8_t *>(e->src_addr);
|
||||
info.des_addr = const_cast<uint8_t *>(e->des_addr);
|
||||
info.rx_ctrl = &rx_ctrl;
|
||||
cb(&info, e->data, static_cast<int>(e->data_len));
|
||||
}
|
||||
|
||||
void on_send(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
// Read the volatile pointer once (see on_recv): disable()/deinit() can null it
|
||||
// on the main loop concurrently with this RX-thread callback.
|
||||
const esp_now_send_cb_t cb = g_send_cb;
|
||||
if (cb == nullptr)
|
||||
return;
|
||||
if (len < sizeof(esp_now_hosted_send_evt_t)) {
|
||||
ESP_LOGW(TAG, "SEND evt too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *e = reinterpret_cast<const esp_now_hosted_send_evt_t *>(data);
|
||||
#if ESP_IDF_VERSION >= ESP_IDF_VERSION_VAL(5, 5, 0)
|
||||
// IDF >= 5.5: esp_now_send_cb_t takes esp_now_send_info_t (== wifi_tx_info_t),
|
||||
// whose des_addr is a POINTER (not an inline array). Point it at the event's
|
||||
// MAC (valid for this callback) — do NOT memcpy into it (that writes NULL and
|
||||
// faults). ESPHome reads only info->des_addr.
|
||||
esp_now_send_info_t si;
|
||||
memset(&si, 0, sizeof(si));
|
||||
si.des_addr = const_cast<uint8_t *>(e->des_addr);
|
||||
cb(&si, static_cast<esp_now_send_status_t>(e->status));
|
||||
#else
|
||||
cb(e->des_addr, static_cast<esp_now_send_status_t>(e->status));
|
||||
#endif
|
||||
}
|
||||
|
||||
esp_err_t ensure_setup() {
|
||||
// Gate on g_setup_done, not on g_req_mutex: a failure part-way through (a
|
||||
// semaphore that did not allocate, a callback that did not register) must not
|
||||
// leave a later call thinking setup completed. Semaphore creation is guarded
|
||||
// so a retry after a partial failure does not leak the earlier handles.
|
||||
if (g_setup_done)
|
||||
return ESP_OK;
|
||||
if (g_req_mutex == nullptr)
|
||||
g_req_mutex = xSemaphoreCreateMutex();
|
||||
if (g_resp_sem == nullptr)
|
||||
g_resp_sem = xSemaphoreCreateBinary();
|
||||
if (g_req_mutex == nullptr || g_resp_sem == nullptr)
|
||||
return ESP_ERR_NO_MEM;
|
||||
esp_err_t err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_RESP, on_resp, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_RECV, on_recv, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_SEND, on_send, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
g_setup_done = true;
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
// Send one request envelope. With wait=true (default) block until the matching
|
||||
// response (or timeout); with wait=false return as soon as the frame is handed
|
||||
// to the transport (fire-and-forget, used by esp_now_send).
|
||||
//
|
||||
// `tail` is an optional second chunk written straight after `payload`. Callers
|
||||
// with a fixed header plus a bulk body (esp_now_send) pass the two separately
|
||||
// so they never need a build buffer of their own: both chunks are laid into the
|
||||
// request buffer here, under g_req_mutex, which keeps concurrent callers from
|
||||
// racing and saves a full copy of the body on every transmit.
|
||||
esp_err_t request(uint8_t opcode, const void *payload, uint16_t plen, void *ret, uint16_t ret_cap, uint16_t *ret_len,
|
||||
bool wait = true, const void *tail = nullptr, uint16_t tail_len = 0) {
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
if (plen > ESP_NOW_HOSTED_MAX_PAYLOAD || tail_len > ESP_NOW_HOSTED_MAX_PAYLOAD - plen)
|
||||
return ESP_ERR_INVALID_SIZE;
|
||||
const uint16_t total_len = static_cast<uint16_t>(plen + tail_len);
|
||||
|
||||
if (xSemaphoreTake(g_req_mutex, portMAX_DELAY) != pdTRUE)
|
||||
return ESP_FAIL;
|
||||
|
||||
static uint8_t buf[sizeof(esp_now_hosted_req_t) + ESP_NOW_HOSTED_MAX_PAYLOAD]; // guarded by g_req_mutex
|
||||
auto *req = reinterpret_cast<esp_now_hosted_req_t *>(buf);
|
||||
req->opcode = opcode;
|
||||
req->seq = ++g_seq;
|
||||
req->payload_len = total_len;
|
||||
if (plen != 0)
|
||||
memcpy(req->payload, payload, plen);
|
||||
if (tail_len != 0)
|
||||
memcpy(req->payload + plen, tail, tail_len);
|
||||
g_expect_seq = req->seq;
|
||||
|
||||
xSemaphoreTake(g_resp_sem, 0); // drain any stale signal before sending
|
||||
err = esp_hosted_send_custom_data(ESP_NOW_HOSTED_MSG_REQ, buf, sizeof(esp_now_hosted_req_t) + total_len);
|
||||
if (err != ESP_OK) {
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return err;
|
||||
}
|
||||
if (!wait) {
|
||||
// Fire-and-forget (esp_now_send): the co-processor enqueues the frame and
|
||||
// reports the real TX result later via the async SEND event, exactly like
|
||||
// native esp_now_send. Returning here keeps the main loop off the ~100 ms+
|
||||
// RPC round-trip. The matching RESP is ignored (seq won't match the next
|
||||
// waited request, so on_resp drops it).
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return ESP_OK;
|
||||
}
|
||||
if (xSemaphoreTake(g_resp_sem, pdMS_TO_TICKS(ESP_NOW_HOSTED_TIMEOUT_MS)) != pdTRUE) {
|
||||
ESP_LOGW(TAG, "opcode %u timed out", opcode);
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return ESP_ERR_TIMEOUT;
|
||||
}
|
||||
|
||||
const int32_t status = g_resp_status;
|
||||
if (ret != nullptr && ret_cap != 0) {
|
||||
uint16_t n = g_resp_ret_len < ret_cap ? g_resp_ret_len : ret_cap;
|
||||
memcpy(ret, const_cast<const uint8_t *>(g_resp_ret), n);
|
||||
if (ret_len != nullptr)
|
||||
*ret_len = n;
|
||||
}
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return static_cast<esp_err_t>(status);
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
// ── The <esp_now.h> surface, defined for the radio-less host ────────────────
|
||||
extern "C" {
|
||||
|
||||
esp_err_t esp_now_init(void) { return request(ESP_NOW_HOSTED_OP_INIT, nullptr, 0, nullptr, 0, nullptr); }
|
||||
|
||||
esp_err_t esp_now_deinit(void) {
|
||||
g_recv_cb = nullptr;
|
||||
g_send_cb = nullptr;
|
||||
peer_cache_clear(); // the co-processor drops all peers on deinit
|
||||
return request(ESP_NOW_HOSTED_OP_DEINIT, nullptr, 0, nullptr, 0, nullptr);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_get_version(uint32_t *version) {
|
||||
uint32_t v = 0;
|
||||
uint16_t rl = 0;
|
||||
esp_err_t err = request(ESP_NOW_HOSTED_OP_GET_VERSION, nullptr, 0, &v, sizeof(v), &rl);
|
||||
if (version != nullptr)
|
||||
*version = v;
|
||||
return err;
|
||||
}
|
||||
|
||||
esp_err_t esp_now_register_recv_cb(esp_now_recv_cb_t cb) {
|
||||
// Only arm the callback once the CustomRpc handlers are actually registered,
|
||||
// so a failed setup leaves g_recv_cb null rather than falsely "registered".
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
g_recv_cb = cb;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_unregister_recv_cb(void) {
|
||||
g_recv_cb = nullptr;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_register_send_cb(esp_now_send_cb_t cb) {
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
g_send_cb = cb;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_unregister_send_cb(void) {
|
||||
g_send_cb = nullptr;
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
static esp_err_t add_or_mod_peer(uint8_t opcode, const esp_now_peer_info_t *peer, bool wait) {
|
||||
if (peer == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
esp_now_hosted_peer_t p;
|
||||
memset(&p, 0, sizeof(p));
|
||||
memcpy(p.peer_addr, peer->peer_addr, 6);
|
||||
memcpy(p.lmk, peer->lmk, 16);
|
||||
p.channel = peer->channel;
|
||||
p.ifidx = static_cast<uint8_t>(peer->ifidx);
|
||||
p.encrypt = peer->encrypt ? 1 : 0;
|
||||
return request(opcode, &p, sizeof(p), nullptr, 0, nullptr, wait);
|
||||
}
|
||||
esp_err_t esp_now_add_peer(const esp_now_peer_info_t *peer) {
|
||||
// Fire-and-forget (wait=false): adding a peer is a blocking RPC round-trip,
|
||||
// and ESPHome's espnow calls it on the main loop when a device joins the mesh
|
||||
// — under co-processor load that stalls the UI (peer-churn stutter). Issue it
|
||||
// without waiting and mirror it locally. Safe against a following
|
||||
// esp_now_send to the same peer: both ride the same in-order CustomRpc
|
||||
// channel (mutex-serialized on the host) and the co-processor processes REQs
|
||||
// FIFO, so ADD_PEER is applied before the SEND. Trade-off: a co-processor-side
|
||||
// failure (e.g. peer table full) is no longer reported synchronously — the
|
||||
// same limitation as esp_now_send — but ESPHome only adds peers it validated.
|
||||
esp_err_t err = add_or_mod_peer(ESP_NOW_HOSTED_OP_ADD_PEER, peer, /*wait=*/false);
|
||||
if (err == ESP_OK)
|
||||
peer_cache_add(peer->peer_addr); // keep the local mirror in sync
|
||||
return err;
|
||||
}
|
||||
esp_err_t esp_now_mod_peer(const esp_now_peer_info_t *peer) {
|
||||
// mod_peer changes a peer's parameters, not its existence, so the cache is
|
||||
// unaffected. Kept synchronous — it is not on any hot path (espnow never
|
||||
// calls it), so the extra round-trip does not matter and the status is useful.
|
||||
return add_or_mod_peer(ESP_NOW_HOSTED_OP_MOD_PEER, peer, /*wait=*/true);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_del_peer(const uint8_t *peer_addr) {
|
||||
if (peer_addr == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
// Fire-and-forget for the same reason as add_peer (peer churn on the main
|
||||
// loop). Removal is order-independent, so this is strictly safe.
|
||||
esp_err_t err = request(ESP_NOW_HOSTED_OP_DEL_PEER, peer_addr, 6, nullptr, 0, nullptr, /*wait=*/false);
|
||||
if (err == ESP_OK)
|
||||
peer_cache_remove(peer_addr); // keep the local mirror in sync
|
||||
return err;
|
||||
}
|
||||
|
||||
bool esp_now_is_peer_exist(const uint8_t *peer_addr) {
|
||||
if (peer_addr == nullptr)
|
||||
return false;
|
||||
// Answered from the local mirror — no RPC round-trip. ESPHome's espnow calls
|
||||
// this on the main loop for every received frame and every send, so a
|
||||
// blocking round-trip here would stall rendering under mesh traffic.
|
||||
return peer_cache_contains(peer_addr);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_send(const uint8_t *peer_addr, const uint8_t *data, size_t len) {
|
||||
if (len > ESP_NOW_HOSTED_MAX_FRAME)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
if (data == nullptr && len != 0) // native esp_now_send treats this as an arg error
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
// Only the small fixed header is built here; the caller's frame goes over as
|
||||
// the request tail, so request() lays both into its own buffer under
|
||||
// g_req_mutex. esp_now_send is a public C symbol and may be called from any
|
||||
// task, and a shared build buffer here would let two callers corrupt each
|
||||
// other's frame. Passing the body through also drops a full-frame copy per
|
||||
// transmit, on the path this shim exists to keep quick.
|
||||
uint8_t hdr[sizeof(esp_now_hosted_send_req_t)];
|
||||
auto *s = reinterpret_cast<esp_now_hosted_send_req_t *>(hdr);
|
||||
s->has_addr = peer_addr != nullptr ? 1 : 0;
|
||||
if (peer_addr != nullptr)
|
||||
memcpy(s->peer_addr, peer_addr, 6);
|
||||
else
|
||||
memset(s->peer_addr, 0, 6);
|
||||
s->data_len = static_cast<uint16_t>(len);
|
||||
// Fire-and-forget (wait=false): native esp_now_send returns once the frame is
|
||||
// queued, with the real TX result delivered later through the send callback.
|
||||
// The co-processor mirrors that — it acks enqueue immediately and reports the
|
||||
// outcome via the async SEND event (on_send -> on_send_report). Waiting for
|
||||
// the RPC RESP here would block the main loop for the full round-trip on
|
||||
// every transmit.
|
||||
return request(ESP_NOW_HOSTED_OP_SEND, hdr, sizeof(hdr), nullptr, 0, nullptr, /*wait=*/false, data,
|
||||
static_cast<uint16_t>(len));
|
||||
}
|
||||
|
||||
esp_err_t esp_now_set_pmk(const uint8_t *pmk) {
|
||||
if (pmk == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
return request(ESP_NOW_HOSTED_OP_SET_PMK, pmk, 16, nullptr, 0, nullptr);
|
||||
}
|
||||
|
||||
// Remainder of the <esp_now.h> surface. Not used by ESPHome's espnow component
|
||||
// today; provided so the whole header links and future callers get a defined
|
||||
// (if unimplemented) symbol rather than a link error. Wire them through
|
||||
// CustomRpc if a use case appears.
|
||||
esp_err_t esp_now_get_peer(const uint8_t * /*peer_addr*/, esp_now_peer_info_t * /*peer*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
esp_err_t esp_now_fetch_peer(bool /*from_head*/, esp_now_peer_info_t * /*peer*/) { return ESP_ERR_NOT_SUPPORTED; }
|
||||
esp_err_t esp_now_get_peer_num(esp_now_peer_num_t * /*num*/) { return ESP_ERR_NOT_SUPPORTED; }
|
||||
esp_err_t esp_now_set_wake_window(uint16_t /*window*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED; // power-save wake window is not forwarded; don't claim success
|
||||
}
|
||||
esp_err_t esp_now_set_peer_rate_config(const uint8_t * /*peer_addr*/, esp_now_rate_config_t * /*cfg*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
esp_err_t esp_wifi_config_espnow_rate(wifi_interface_t /*ifx*/, wifi_phy_rate_t /*rate*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
|
||||
} // extern "C"
|
||||
|
||||
#endif // CONFIG_IDF_TARGET_ESP32P4
|
||||
@@ -0,0 +1,128 @@
|
||||
/*
|
||||
* esp_now_hosted — ESP-NOW-over-CustomRpc wire protocol.
|
||||
*
|
||||
* Shared, byte-for-byte-identical contract between:
|
||||
* - the host shim (esphome/components/esp32_hosted/esp_now_hosted.cpp)
|
||||
* - the coprocessor firmware (esphome/esp-hosted-firmware)
|
||||
*
|
||||
* It rides esp-hosted's CustomRpc channel (RPC ID 388, "peer data transfer",
|
||||
* available since esp-hosted v2.8.1), teaching the radio-less host <-> radio
|
||||
* co-processor link to carry esp_now.h, which esp-hosted itself does not proxy
|
||||
* (Espressif issue espressif/esp-hosted-mcu#19).
|
||||
*
|
||||
* KEEP THE TWO COPIES IN SYNC. The canonical copy lives here; the coprocessor
|
||||
* firmware uses a verbatim copy. Both sides are little-endian, so these packed
|
||||
* structs are wire-compatible with no byte-swapping.
|
||||
*/
|
||||
|
||||
#ifndef ESP_NOW_HOSTED_RPC_H
|
||||
#define ESP_NOW_HOSTED_RPC_H
|
||||
|
||||
#ifdef __cplusplus
|
||||
#include <cstdint>
|
||||
#else
|
||||
#include <stdint.h>
|
||||
#endif
|
||||
|
||||
#ifdef __cplusplus
|
||||
extern "C" {
|
||||
#endif
|
||||
|
||||
/* ── CustomRpc message IDs (any uint32_t except 0xFFFFFFFF) ──────────────────
|
||||
* One REQ handler slot on the device; three event handler slots on the host.
|
||||
* The bytes spell "now" + index, a private range unlikely to clash with other
|
||||
* CustomRpc users (e.g. the stock peer_data_transfer example's 1..6). */
|
||||
#define ESP_NOW_HOSTED_MSG_REQ 0x6E6F7701u /* host -> device : request envelope */
|
||||
#define ESP_NOW_HOSTED_MSG_RESP 0x6E6F7702u /* device -> host : reply to a REQ */
|
||||
#define ESP_NOW_HOSTED_MSG_RECV 0x6E6F7703u /* device -> host : async RX frame */
|
||||
#define ESP_NOW_HOSTED_MSG_SEND 0x6E6F7704u /* device -> host : async TX status */
|
||||
|
||||
/* ── Request opcodes ────────────────────────────────────────────────────── */
|
||||
enum {
|
||||
ESP_NOW_HOSTED_OP_INIT = 1, /* esp_now_init + register device recv/send cbs */
|
||||
ESP_NOW_HOSTED_OP_DEINIT = 2, /* unregister cbs + esp_now_deinit */
|
||||
ESP_NOW_HOSTED_OP_ADD_PEER = 3, /* payload: esp_now_hosted_peer_t */
|
||||
ESP_NOW_HOSTED_OP_DEL_PEER = 4, /* payload: 6-byte peer MAC */
|
||||
ESP_NOW_HOSTED_OP_IS_PEER_EXIST = 5, /* payload: 6-byte MAC; ret: 1 byte bool */
|
||||
ESP_NOW_HOSTED_OP_SEND = 6, /* payload: esp_now_hosted_send_req_t */
|
||||
ESP_NOW_HOSTED_OP_GET_VERSION = 7, /* ret: uint32 version */
|
||||
ESP_NOW_HOSTED_OP_SET_PMK = 8, /* payload: 16-byte PMK */
|
||||
ESP_NOW_HOSTED_OP_MOD_PEER = 9, /* payload: esp_now_hosted_peer_t */
|
||||
};
|
||||
|
||||
/* Largest ESP-NOW payload we forward. ESP-NOW v2 (IDF >= 5.4) is 1470 B; well
|
||||
* under esp-hosted's 8166 B CustomRpc cap, so the shim never truncates. */
|
||||
#define ESP_NOW_HOSTED_MAX_FRAME 1470u
|
||||
/* Envelope slack for the largest opcode payload (a SEND req wrapping a frame). */
|
||||
#define ESP_NOW_HOSTED_MAX_PAYLOAD (ESP_NOW_HOSTED_MAX_FRAME + 16u)
|
||||
/* Host request/response round-trip timeout over the transport. Generous:
|
||||
* normal RTT is sub-millisecond, but Wi-Fi/BLE contention on the co-processor
|
||||
* can stall the RX thread. */
|
||||
#define ESP_NOW_HOSTED_TIMEOUT_MS 2000
|
||||
|
||||
/* ── Envelopes ──────────────────────────────────────────────────────────── */
|
||||
|
||||
/* These payloads are shared verbatim with the C co-processor firmware, so they
|
||||
* use C's `typedef struct {...} name;` idiom rather than C++ `using` aliases,
|
||||
* which would not compile there. Silence clang-tidy's modernize-use-using for
|
||||
* the shared struct block. */
|
||||
// NOLINTBEGIN(modernize-use-using)
|
||||
typedef struct {
|
||||
uint8_t opcode; /* one of ESP_NOW_HOSTED_OP_* */
|
||||
uint8_t seq; /* wraps 0..255; echoed in the response for matching */
|
||||
uint16_t payload_len; /* bytes of opcode-specific payload that follow */
|
||||
uint8_t payload[]; /* flexible */
|
||||
} __attribute__((packed)) esp_now_hosted_req_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t opcode; /* echoes the request opcode */
|
||||
uint8_t seq; /* echoes the request seq */
|
||||
int32_t status; /* esp_err_t from the native call on the co-processor */
|
||||
uint16_t ret_len; /* bytes of return payload that follow */
|
||||
uint8_t ret[]; /* flexible (e.g. version u32, is_peer_exist bool) */
|
||||
} __attribute__((packed)) esp_now_hosted_resp_t;
|
||||
|
||||
/* ── Opcode payloads ────────────────────────────────────────────────────── */
|
||||
|
||||
/* esp_now_peer_info_t minus the host-only `priv` pointer, which is meaningless
|
||||
* across the transport and never set by ESPHome's espnow component. */
|
||||
typedef struct {
|
||||
uint8_t peer_addr[6];
|
||||
uint8_t lmk[16];
|
||||
uint8_t channel; /* 0 = current channel */
|
||||
uint8_t ifidx; /* wifi_interface_t (0=STA, 1=AP) */
|
||||
uint8_t encrypt; /* bool */
|
||||
} __attribute__((packed)) esp_now_hosted_peer_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t has_addr; /* 0 => peer_addr is NULL (broadcast to all peers) */
|
||||
uint8_t peer_addr[6];
|
||||
uint16_t data_len;
|
||||
uint8_t data[]; /* flexible, up to ESP_NOW_HOSTED_MAX_FRAME */
|
||||
} __attribute__((packed)) esp_now_hosted_send_req_t;
|
||||
|
||||
/* ── Async events (device -> host) ──────────────────────────────────────── */
|
||||
|
||||
/* Reconstructed on the host into an esp_now_recv_info_t + a minimal
|
||||
* wifi_pkt_rx_ctrl_t. ESPHome's espnow reads info->src_addr, info->des_addr,
|
||||
* info->rx_ctrl->rssi and info->rx_ctrl->timestamp. */
|
||||
typedef struct {
|
||||
uint8_t src_addr[6];
|
||||
uint8_t des_addr[6];
|
||||
int8_t rssi;
|
||||
uint8_t channel;
|
||||
uint16_t data_len;
|
||||
uint8_t data[]; /* flexible */
|
||||
} __attribute__((packed)) esp_now_hosted_recv_evt_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t des_addr[6];
|
||||
uint8_t status; /* esp_now_send_status_t (0 = success) */
|
||||
} __attribute__((packed)) esp_now_hosted_send_evt_t;
|
||||
// NOLINTEND(modernize-use-using)
|
||||
|
||||
#ifdef __cplusplus
|
||||
}
|
||||
#endif
|
||||
|
||||
#endif /* ESP_NOW_HOSTED_RPC_H */
|
||||
@@ -3,6 +3,7 @@ from typing import Any
|
||||
from esphome import automation, core
|
||||
import esphome.codegen as cg
|
||||
from esphome.components import wifi
|
||||
from esphome.components.esp32 import VARIANT_ESP32P4, get_esp32_variant
|
||||
from esphome.components.udp import CONF_ON_RECEIVE
|
||||
import esphome.config_validation as cv
|
||||
from esphome.const import (
|
||||
@@ -17,6 +18,7 @@ from esphome.const import (
|
||||
)
|
||||
from esphome.core import CORE, HexInt
|
||||
from esphome.cpp_generator import MockObj, TemplateArgsType
|
||||
import esphome.final_validate as fv
|
||||
from esphome.types import ConfigType
|
||||
|
||||
CODEOWNERS = ["@jesserockz"]
|
||||
@@ -132,6 +134,24 @@ CONFIG_SCHEMA = cv.All(
|
||||
)
|
||||
|
||||
|
||||
def _validate_variant(config: ConfigType) -> ConfigType:
|
||||
# ESP-NOW rides the Wi-Fi PHY. Radio-less esp32 variants have no native
|
||||
# ESP-NOW; only the ESP32-P4 has a path, via the esp32_hosted shim that
|
||||
# supplies the esp_now_* symbols. Fail here with a clear message instead of
|
||||
# letting the build reach an "undefined reference to esp_now_*" link error.
|
||||
variant = get_esp32_variant()
|
||||
if wifi.variant_has_wifi(variant):
|
||||
return config
|
||||
if variant != VARIANT_ESP32P4:
|
||||
raise cv.Invalid(f"ESP-NOW is not supported on {variant} (no Wi-Fi radio)")
|
||||
if "esp32_hosted" not in fv.full_config.get():
|
||||
raise cv.Invalid(f"ESP-NOW on {variant} requires the esp32_hosted component")
|
||||
return config
|
||||
|
||||
|
||||
FINAL_VALIDATE_SCHEMA = _validate_variant
|
||||
|
||||
|
||||
async def _trigger_to_code(config: ConfigType) -> MockObj:
|
||||
if address := config.get(CONF_ADDRESS):
|
||||
address = address.parts
|
||||
|
||||
@@ -192,6 +192,8 @@ async def to_code(config: ConfigType) -> None:
|
||||
if CORE.using_arduino:
|
||||
if CORE.is_esp8266:
|
||||
cg.add_library("ESP8266mDNS", None)
|
||||
# No MDNS global in the build; mdns_esp8266.cpp owns a guarded MDNSResponder
|
||||
cg.add_build_flag("-DNO_GLOBAL_MDNS")
|
||||
elif CORE.is_rp2:
|
||||
cg.add_library("LEAmDNS", None)
|
||||
|
||||
|
||||
@@ -13,8 +13,47 @@
|
||||
|
||||
namespace esphome::mdns {
|
||||
|
||||
// Main-loop calls into LEAmDNS that send (update() and close(); begin(), addService() and
|
||||
// the scheduled restart never reach a send) can yield inside UdpContext::sendTimeout(); a
|
||||
// packet arriving then re-enters LEAmDNS from lwIP on the same UdpContext and both sides
|
||||
// free the same tx pbufs (#18760). Received packets stay queued during such a call and are
|
||||
// processed from the main loop afterwards.
|
||||
class GuardedMDNSResponder : public ::esp8266::MDNSImplementation::MDNSResponder {
|
||||
public:
|
||||
void update_guarded() { this->run_guarded_(&GuardedMDNSResponder::update); }
|
||||
void close_guarded() { this->run_guarded_(&GuardedMDNSResponder::close); }
|
||||
|
||||
private:
|
||||
void run_guarded_(bool (GuardedMDNSResponder::*fn)()) {
|
||||
UdpContext *ctx = this->m_pUDPContext;
|
||||
if (ctx == nullptr) {
|
||||
(this->*fn)();
|
||||
return;
|
||||
}
|
||||
// Set every time: a restart replaces the context together with its stock handler. Only
|
||||
// begin() and the scheduled netif callback restart, never update() or close(), so the
|
||||
// context cannot change underneath this call.
|
||||
ctx->onRx([this]() {
|
||||
if (!this->in_loop_call_) {
|
||||
this->_callProcess();
|
||||
}
|
||||
});
|
||||
this->in_loop_call_ = true;
|
||||
(this->*fn)();
|
||||
// close() releases the context; a yield in here queues further packets for this loop too
|
||||
while (this->m_pUDPContext != nullptr && this->m_pUDPContext->next()) {
|
||||
this->_parseMessage();
|
||||
}
|
||||
this->in_loop_call_ = false;
|
||||
}
|
||||
|
||||
volatile bool in_loop_call_{false};
|
||||
};
|
||||
|
||||
static GuardedMDNSResponder mdns_responder; // NOLINT(cppcoreguidelines-avoid-non-const-global-variables)
|
||||
|
||||
static void register_esp8266(MDNSComponent *, StaticVector<MDNSService, MDNS_SERVICE_COUNT> &services) {
|
||||
MDNS.begin(App.get_name().c_str());
|
||||
mdns_responder.begin(App.get_name().c_str());
|
||||
|
||||
for (const auto &service : services) {
|
||||
// Strip the leading underscore from the proto and service_type. While it is
|
||||
@@ -30,10 +69,10 @@ static void register_esp8266(MDNSComponent *, StaticVector<MDNSService, MDNS_SER
|
||||
service_type++;
|
||||
}
|
||||
uint16_t port = service.port.value();
|
||||
MDNS.addService(FPSTR(service_type), FPSTR(proto), port);
|
||||
mdns_responder.addService(FPSTR(service_type), FPSTR(proto), port);
|
||||
for (const auto &record : service.txt_records) {
|
||||
MDNS.addServiceTxt(FPSTR(service_type), FPSTR(proto), FPSTR(MDNS_STR_ARG(record.key)),
|
||||
FPSTR(MDNS_STR_ARG(record.value)));
|
||||
mdns_responder.addServiceTxt(FPSTR(service_type), FPSTR(proto), FPSTR(MDNS_STR_ARG(record.key)),
|
||||
FPSTR(MDNS_STR_ARG(record.value)));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -52,7 +91,7 @@ void MDNSComponent::start_polling_window_() {
|
||||
if (wifi->is_roaming() || (!wifi->is_connected() && !wifi->is_ap_active()))
|
||||
return;
|
||||
#endif
|
||||
MDNS.update();
|
||||
mdns_responder.update_guarded();
|
||||
});
|
||||
this->set_timeout(MDNS_POLL_STOP_ID, MDNS_POLL_WINDOW_MS, [this]() { this->cancel_interval(MDNS_POLL_ID); });
|
||||
}
|
||||
@@ -81,7 +120,7 @@ void MDNSComponent::on_ip_state(const network::IPAddresses &ips, const network::
|
||||
#endif
|
||||
|
||||
void MDNSComponent::on_shutdown() {
|
||||
MDNS.close();
|
||||
mdns_responder.close_guarded();
|
||||
delay(10);
|
||||
}
|
||||
|
||||
|
||||
@@ -71,6 +71,7 @@
|
||||
#define USE_ESP32_HOSTED
|
||||
#define USE_ESP32_HOSTED_HTTP_UPDATE
|
||||
#define USE_ESP32_IMPROV_STATE_CALLBACK
|
||||
#define USE_ESP_NOW_HOSTED
|
||||
#define USE_EVENT
|
||||
#define USE_FAN
|
||||
#define USE_GPIO_BINARY_SENSOR_INTERRUPT
|
||||
|
||||
@@ -23,6 +23,7 @@ from esphome.net_retry import (
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from filelock import FileLock
|
||||
import requests
|
||||
|
||||
PathType = str | os.PathLike
|
||||
@@ -909,6 +910,61 @@ def _part_path(dest: Path) -> Path:
|
||||
return dest.with_name(dest.name + ".part")
|
||||
|
||||
|
||||
def downloaded_bytes(dest: Path, size: int | None = None) -> int:
|
||||
"""Bytes of ``dest`` on disk (its ``.part`` while streaming), capped at ``size``."""
|
||||
done = 0
|
||||
for candidate in (_part_path(dest), dest):
|
||||
try:
|
||||
done = candidate.stat().st_size
|
||||
break
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
return done if size is None else min(done, size)
|
||||
|
||||
|
||||
# Short lock-acquire slices so a waiting worker still observes Ctrl-C
|
||||
_DOWNLOAD_LOCK_POLL = 1
|
||||
|
||||
# Waiting on another process's download; past this the caller leaves the
|
||||
# file to its holder (the later sequential install waits on the same lock)
|
||||
DOWNLOAD_LOCK_TIMEOUT = 60
|
||||
|
||||
|
||||
class DownloadLockUnavailable(OSError):
|
||||
"""The lock file cannot be used at all (a lock-less filesystem)."""
|
||||
|
||||
|
||||
def wait_for_download_lock(
|
||||
lock: "FileLock",
|
||||
tracker: Callable[[int], None],
|
||||
on_disk: Callable[[], int],
|
||||
name: str,
|
||||
) -> None:
|
||||
"""Acquire ``lock``, reporting ``on_disk()`` to ``tracker`` each poll so the
|
||||
bar follows the holder's download. Raises filelock's ``Timeout`` once
|
||||
``DOWNLOAD_LOCK_TIMEOUT`` seconds pass."""
|
||||
from filelock import Timeout
|
||||
|
||||
deadline = time.monotonic() + DOWNLOAD_LOCK_TIMEOUT
|
||||
waiting = False
|
||||
while True:
|
||||
try:
|
||||
lock.acquire(timeout=_DOWNLOAD_LOCK_POLL)
|
||||
return
|
||||
except Timeout:
|
||||
pass
|
||||
except OSError as err:
|
||||
# Distinct from an OSError out of on_disk(), which must not
|
||||
# read as "locks unsupported"
|
||||
raise DownloadLockUnavailable(*err.args) from err
|
||||
if not waiting:
|
||||
waiting = True
|
||||
_LOGGER.info("Waiting for another process downloading %s", name)
|
||||
tracker(on_disk()) # raises when the batch is cancelled
|
||||
if time.monotonic() >= deadline:
|
||||
raise Timeout(lock.lock_file)
|
||||
|
||||
|
||||
def discard_partial_download(dest: Path) -> None:
|
||||
"""Remove ``dest`` and the resume sidecars of an abandoned download."""
|
||||
part = _part_path(dest)
|
||||
@@ -1319,10 +1375,7 @@ def download_from_mirrors(
|
||||
)
|
||||
# Tick with the bytes already on disk so a combined bar holds
|
||||
# steady during the backoff instead of rewinding to zero
|
||||
done = 0
|
||||
if progress is not None:
|
||||
part = _part_path(path_target)
|
||||
done = part.stat().st_size if part.is_file() else 0
|
||||
done = downloaded_bytes(path_target) if progress is not None else 0
|
||||
_cancellable_sleep(delay, progress, done)
|
||||
|
||||
# 3. Report every attempted URL if all mirrors failed. failures spans
|
||||
|
||||
@@ -33,11 +33,14 @@ import time
|
||||
from typing import Any, NamedTuple
|
||||
|
||||
from esphome.framework_helpers import (
|
||||
DownloadLockUnavailable,
|
||||
content_length,
|
||||
discard_partial_download,
|
||||
downloaded_bytes,
|
||||
failure_reason,
|
||||
resume_fetch_job,
|
||||
run_batch_downloads,
|
||||
wait_for_download_lock,
|
||||
warn_prefetch_failures,
|
||||
)
|
||||
from esphome.helpers import get_bool_env, get_usable_cpu_count, rmtree
|
||||
@@ -61,16 +64,10 @@ _RESOLVE_WORKERS = 8
|
||||
# A hung child must not block the build; downloads resume on the next run
|
||||
_PREFETCH_TIMEOUT = 20 * 60
|
||||
|
||||
# Waiting on another process's URL download; past this, leave it to pio
|
||||
_DOWNLOAD_LOCK_TIMEOUT = 60
|
||||
|
||||
# Child exit for a handled, already-warned failure; 1 would collide with
|
||||
# the interpreter's own import-failure exit
|
||||
_EXIT_HANDLED = 3
|
||||
|
||||
# Short lock-acquire slices so a waiting worker still observes Ctrl-C
|
||||
_URI_LOCK_POLL = 1
|
||||
|
||||
# Resolution errored (vs a clean skip); suppresses the warm sentinel
|
||||
_RESOLVE_FAILED = object()
|
||||
|
||||
@@ -462,51 +459,54 @@ def _uri_jobs(
|
||||
|
||||
|
||||
def _serialized_fetch_job(
|
||||
dl_path: Path, lock_path: str, body: Any, unlocked_ok: bool = True
|
||||
dl_path: Path,
|
||||
lock_path: str,
|
||||
body: Any,
|
||||
size: int,
|
||||
stream_dest: Path | None = None,
|
||||
unlocked_ok: bool = True,
|
||||
) -> Any:
|
||||
"""Wrap ``body`` so the shared destination is single-writer.
|
||||
|
||||
Interleaved writers truncate each other's ``.part`` bytes (see
|
||||
registry.py). The bounded poll observes Ctrl-C via the tracker; a
|
||||
blown deadline is a clean skip (the holder's copy is what the build
|
||||
needs). On a lock-less filesystem a sha256-verified body runs
|
||||
unlocked with one warning; a checksum-less one
|
||||
(``unlocked_ok=False``) is a counted failure instead.
|
||||
"""Wrap ``body`` so the shared destination is single-writer (interleaved
|
||||
writers truncate each other's ``.part``, see registry.py). A blown deadline
|
||||
is a clean skip. On a lock-less filesystem a sha256-verified body runs
|
||||
unlocked with one warning; a checksum-less one (``unlocked_ok=False``) fails.
|
||||
"""
|
||||
|
||||
def on_disk() -> int:
|
||||
# A URL job's holder streams beside the staging path until it
|
||||
# promotes; after that only dl_path is left
|
||||
done = downloaded_bytes(dl_path, size)
|
||||
if not done and stream_dest is not None:
|
||||
done = downloaded_bytes(stream_dest, size)
|
||||
return done
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
from filelock import FileLock, Timeout
|
||||
|
||||
# fallback_to_soft would leave a stale marker on lock-less
|
||||
# filesystems that blocks every later build (see git.py)
|
||||
lock = FileLock(lock_path, fallback_to_soft=False)
|
||||
deadline = time.monotonic() + _DOWNLOAD_LOCK_TIMEOUT
|
||||
while True:
|
||||
try:
|
||||
lock.acquire(timeout=_URI_LOCK_POLL)
|
||||
break
|
||||
except Timeout:
|
||||
tracker(0) # raises when the batch is cancelled
|
||||
if time.monotonic() >= deadline:
|
||||
# Another process is fetching this same file; its copy
|
||||
# is what the build needs (a large framework archive
|
||||
# can hold the lock far longer than this deadline)
|
||||
_LOGGER.debug("Leaving %s to its current downloader", dl_path.name)
|
||||
return
|
||||
except OSError as err:
|
||||
if not unlocked_ok:
|
||||
# A body with no checksum to catch interleaved corruption
|
||||
raise
|
||||
lock = None
|
||||
_LOGGER.warning(
|
||||
"Could not lock %s (%s); downloading unlocked",
|
||||
dl_path.name,
|
||||
err,
|
||||
)
|
||||
break
|
||||
try:
|
||||
wait_for_download_lock(lock, tracker, on_disk, dl_path.name)
|
||||
except Timeout:
|
||||
# The holder's copy is what the build needs (a large
|
||||
# framework archive can outlast this deadline)
|
||||
_LOGGER.debug("Leaving %s to its current downloader", dl_path.name)
|
||||
return
|
||||
except DownloadLockUnavailable as err:
|
||||
if not unlocked_ok:
|
||||
# A body with no checksum to catch interleaved corruption
|
||||
raise
|
||||
lock = None
|
||||
_LOGGER.warning(
|
||||
"Could not lock %s (%s); downloading unlocked",
|
||||
dl_path.name,
|
||||
err,
|
||||
)
|
||||
try:
|
||||
if dl_path.is_file():
|
||||
return # another process finished it while we waited
|
||||
tracker(size) # another process finished it while we waited
|
||||
return
|
||||
body(tracker)
|
||||
finally:
|
||||
if lock is not None:
|
||||
@@ -540,6 +540,7 @@ def _registry_fetch_job(
|
||||
dl_path,
|
||||
f"{dl_path}.esphome.lock",
|
||||
resume_fetch_job(url, dl_path, sha256=checksum, size=size),
|
||||
size,
|
||||
)
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
@@ -571,9 +572,9 @@ def _uri_fetch_job(manager: Any, url: str, dl_path: Path, size: int) -> Any:
|
||||
tmp.replace(dl_path)
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
_serialized_fetch_job(dl_path, f"{tmp}.lock", promote, unlocked_ok=False)(
|
||||
tracker
|
||||
)
|
||||
_serialized_fetch_job(
|
||||
dl_path, f"{tmp}.lock", promote, size, tmp, unlocked_ok=False
|
||||
)(tracker)
|
||||
if dl_path.is_file():
|
||||
# Won or lost, the race is over; staging files left behind
|
||||
# are dead weight PlatformIO's cache never prunes
|
||||
|
||||
@@ -17,8 +17,10 @@ from esphome.framework_helpers import (
|
||||
archive_extract_all,
|
||||
download_from_mirrors,
|
||||
download_with_resume,
|
||||
downloaded_bytes,
|
||||
rmdir,
|
||||
run_batch_downloads,
|
||||
wait_for_download_lock,
|
||||
)
|
||||
from esphome.net_retry import fetch_with_retry, http_request
|
||||
|
||||
@@ -164,11 +166,17 @@ class _PendingArchive(NamedTuple):
|
||||
name: str
|
||||
version: str
|
||||
dest: Path
|
||||
archive: Path
|
||||
url: str
|
||||
sha256: str
|
||||
size: int
|
||||
|
||||
|
||||
def _archive_path(downloads_dir: Path, name: str, version: str) -> Path:
|
||||
"""The one archive path the prefetch and the sequential install share."""
|
||||
return downloads_dir / f"{name}-{version}"
|
||||
|
||||
|
||||
def _already_installed(dest: Path) -> bool:
|
||||
"""Whether ``dest`` holds a completed install (extraction marker)."""
|
||||
return (dest / ".esphome_extracted").is_file()
|
||||
@@ -187,18 +195,18 @@ def prefetch_packages(
|
||||
lock as ``install_package``: the archive's ``.part`` file is shared, and
|
||||
two concurrent writers would truncate each other's bytes.
|
||||
"""
|
||||
from filelock import FileLock
|
||||
from filelock import FileLock, Timeout
|
||||
|
||||
pending: list[_PendingArchive] = []
|
||||
seen: set[str] = set()
|
||||
seen: set[Path] = set()
|
||||
for name, version, dest, mirrors in packages:
|
||||
if mirrors or (dest / ".esphome_extracted").is_file():
|
||||
continue
|
||||
archive_name = f"{name}-{version}"
|
||||
if archive_name in seen:
|
||||
archive = _archive_path(downloads_dir, name, version)
|
||||
if archive in seen:
|
||||
# A duplicate entry would race itself between two workers
|
||||
continue
|
||||
seen.add(archive_name)
|
||||
seen.add(archive)
|
||||
try:
|
||||
url, sha256, size = registry_download(name, version)
|
||||
except EsphomeError as err:
|
||||
@@ -207,10 +215,9 @@ def prefetch_packages(
|
||||
continue
|
||||
if not size:
|
||||
continue
|
||||
archive = downloads_dir / archive_name
|
||||
if archive.is_file() and archive.stat().st_size == size:
|
||||
continue
|
||||
pending.append(_PendingArchive(name, version, dest, url, sha256, size))
|
||||
pending.append(_PendingArchive(name, version, dest, archive, url, sha256, size))
|
||||
if len(pending) < 2:
|
||||
return
|
||||
downloads_dir.mkdir(parents=True, exist_ok=True)
|
||||
@@ -222,20 +229,36 @@ def prefetch_packages(
|
||||
|
||||
def _fetch(entry: _PendingArchive, tracker: Callable[[int], None]) -> None:
|
||||
entry.dest.parent.mkdir(parents=True, exist_ok=True)
|
||||
with FileLock(f"{entry.dest}.lock", fallback_to_soft=False):
|
||||
# Marker re-check: a concurrent build may have installed (and
|
||||
# deleted the archive of) this package while we waited;
|
||||
# re-downloading would orphan a fresh copy in downloads_dir
|
||||
# no branch: the thread tracer misses the skip edge; both
|
||||
# arms of _already_installed are pinned directly
|
||||
if not _already_installed(entry.dest): # pragma: no branch
|
||||
download_with_resume(
|
||||
entry.url,
|
||||
downloads_dir / f"{entry.name}-{entry.version}",
|
||||
sha256=entry.sha256,
|
||||
size=entry.size,
|
||||
progress=tracker,
|
||||
)
|
||||
|
||||
def on_disk() -> int:
|
||||
if done := downloaded_bytes(entry.archive, entry.size):
|
||||
return done
|
||||
# The holder deletes the archive once it has installed it
|
||||
return entry.size if _already_installed(entry.dest) else 0
|
||||
|
||||
lock = FileLock(f"{entry.dest}.lock", fallback_to_soft=False)
|
||||
try:
|
||||
wait_for_download_lock(lock, tracker, on_disk, entry.name)
|
||||
except Timeout:
|
||||
# install_package waits on this same lock and verifies the
|
||||
# holder's copy
|
||||
_LOGGER.debug("Leaving %s to its current downloader", entry.name)
|
||||
return
|
||||
try:
|
||||
if _already_installed(entry.dest):
|
||||
# A concurrent build installed it while we waited; a
|
||||
# re-download would orphan a fresh copy in downloads_dir
|
||||
tracker(entry.size)
|
||||
return
|
||||
download_with_resume(
|
||||
entry.url,
|
||||
entry.archive,
|
||||
sha256=entry.sha256,
|
||||
size=entry.size,
|
||||
progress=tracker,
|
||||
)
|
||||
finally:
|
||||
lock.release()
|
||||
|
||||
failures = run_batch_downloads(
|
||||
"Downloading packages",
|
||||
@@ -288,7 +311,7 @@ def install_package(
|
||||
rmdir(dest, msg=f"Clean up incomplete {name} install")
|
||||
# Persistent location so an interrupted download resumes across runs.
|
||||
downloads_dir.mkdir(parents=True, exist_ok=True)
|
||||
archive = downloads_dir / f"{name}-{version}"
|
||||
archive = _archive_path(downloads_dir, name, version)
|
||||
_LOGGER.info("Downloading %s %s ...", name, version)
|
||||
if mirrors:
|
||||
_LOGGER.warning(
|
||||
|
||||
+16
-1
@@ -294,6 +294,9 @@ def highlight(s):
|
||||
"esphome/components/socket/headers.h",
|
||||
"esphome/core/defines.h",
|
||||
"esphome/components/http_request/httplib.h",
|
||||
# Shared C wire header (byte-identical with the co-processor firmware);
|
||||
# these are protocol constants and constexpr is C++-only.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
def lint_no_defines(fname, match):
|
||||
@@ -816,6 +819,10 @@ def lint_relative_py_import(fname: Path, line, col, content):
|
||||
"esphome/components/host/helpers.cpp",
|
||||
"esphome/components/zephyr/helpers.cpp",
|
||||
"esphome/components/http_request/httplib.h",
|
||||
# Global extern "C" esp_now_* linker symbols + shared C wire header;
|
||||
# neither can live in a C++ namespace.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted.cpp",
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
def lint_namespace(fname: Path, content: str) -> str | None:
|
||||
@@ -841,7 +848,15 @@ def lint_esphome_h(fname, line, col, content):
|
||||
)
|
||||
|
||||
|
||||
@lint_content_check(include=["*.h"], exclude=["esphome/core/entity_types.h"])
|
||||
@lint_content_check(
|
||||
include=["*.h"],
|
||||
exclude=[
|
||||
"esphome/core/entity_types.h",
|
||||
# Shared C wire header; uses a classic #ifndef guard for portability
|
||||
# across the co-processor firmware repo it stays byte-identical with.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
def lint_pragma_once(fname, content):
|
||||
if "#pragma once" not in content:
|
||||
return (
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
# Exercises the ESP-NOW-over-hosted shim: on the ESP32-P4 host, esp32_hosted
|
||||
# supplies the esp_now_* symbols that the espnow component links against.
|
||||
packages:
|
||||
esp32_hosted: !include common.yaml
|
||||
espnow: !include ../espnow/common.yaml
|
||||
@@ -1,143 +1,154 @@
|
||||
{
|
||||
"tests/integration/test_action_concurrent_reentry.py": 57.91,
|
||||
"tests/integration/test_addressable_light_transition.py": 21.25,
|
||||
"tests/integration/test_alarm_control_panel_state_transitions.py": 70.71,
|
||||
"tests/integration/test_api_action_metadata.py": 66.6,
|
||||
"tests/integration/test_api_action_responses.py": 36.1,
|
||||
"tests/integration/test_api_action_timeout.py": 68.86,
|
||||
"tests/integration/test_api_conditional_memory.py": 15.48,
|
||||
"tests/integration/test_api_custom_services.py": 18.77,
|
||||
"tests/integration/test_api_get_time_response_timezone.py": 21.08,
|
||||
"tests/integration/test_api_homeassistant.py": 65.59,
|
||||
"tests/integration/test_api_homeassistant_action_no_subscriber.py": 18.44,
|
||||
"tests/integration/test_api_homeassistant_binary_sensor_initial_state.py": 15.05,
|
||||
"tests/integration/test_api_list_entities_backpressure.py": 13.88,
|
||||
"tests/integration/test_api_message_size_batching.py": 29.98,
|
||||
"tests/integration/test_api_reboot_timeout.py": 16.05,
|
||||
"tests/integration/test_api_string_lambda.py": 15.31,
|
||||
"tests/integration/test_api_vv_logging.py": 19.28,
|
||||
"tests/integration/test_api_zero_psk_provisioning.py": 31.5,
|
||||
"tests/integration/test_areas_and_devices.py": 24.95,
|
||||
"tests/integration/test_automation_wait_actions.py": 20.92,
|
||||
"tests/integration/test_automations.py": 35.19,
|
||||
"tests/integration/test_batch_delay_zero_rapid_transitions.py": 17.99,
|
||||
"tests/integration/test_binary_sensor_autorepeat_filter.py": 20.39,
|
||||
"tests/integration/test_binary_sensor_invalidate_state.py": 18.41,
|
||||
"tests/integration/test_blocking_warning_log_time_not_charged_to_next_operation.py": 24.69,
|
||||
"tests/integration/test_build_info.py": 18.7,
|
||||
"tests/integration/test_camera_mock.py": 16.23,
|
||||
"tests/integration/test_climate_control_action.py": 21.14,
|
||||
"tests/integration/test_climate_custom_modes.py": 20.74,
|
||||
"tests/integration/test_continuation_actions.py": 16.81,
|
||||
"tests/integration/test_cover_control_action.py": 20.34,
|
||||
"tests/integration/test_crc8_helper.py": 9.36,
|
||||
"tests/integration/test_device_id_in_state.py": 44.67,
|
||||
"tests/integration/test_duplicate_entities.py": 23.58,
|
||||
"tests/integration/test_entity_icon.py": 34.35,
|
||||
"tests/integration/test_fan_turn_on_action.py": 24.23,
|
||||
"tests/integration/test_fnv1_hash_object_id.py": 16.21,
|
||||
"tests/integration/test_fnv1a_hash.py": 13.38,
|
||||
"tests/integration/test_gpio_expander_cache.py": 13.06,
|
||||
"tests/integration/test_host_logger_thread_safety.py": 23.66,
|
||||
"tests/integration/test_host_mode_basic.py": 8.01,
|
||||
"tests/integration/test_host_mode_batch_delay.py": 21.0,
|
||||
"tests/integration/test_host_mode_climate_basic_state.py": 22.14,
|
||||
"tests/integration/test_host_mode_climate_control.py": 19.39,
|
||||
"tests/integration/test_host_mode_empty_string_options.py": 21.76,
|
||||
"tests/integration/test_host_mode_entity_fields.py": 29.61,
|
||||
"tests/integration/test_host_mode_fan_preset.py": 20.01,
|
||||
"tests/integration/test_host_mode_many_entities.py": 39.08,
|
||||
"tests/integration/test_host_mode_many_entities_multiple_connections.py": 23.92,
|
||||
"tests/integration/test_host_mode_noise_encryption.py": 42.42,
|
||||
"tests/integration/test_host_mode_reconnect.py": 3.41,
|
||||
"tests/integration/test_host_mode_sensor.py": 22.96,
|
||||
"tests/integration/test_host_ota.py": 29.5,
|
||||
"tests/integration/test_host_preferences.py": 16.06,
|
||||
"tests/integration/test_host_preferences_suspend_resume.py": 18.71,
|
||||
"tests/integration/test_improv_serial_uart.py": 20.22,
|
||||
"tests/integration/test_large_message_batching.py": 26.56,
|
||||
"tests/integration/test_legacy_area.py": 22.72,
|
||||
"tests/integration/test_legacy_climate_compat.py": 14.13,
|
||||
"tests/integration/test_legacy_fan_compat.py": 14.33,
|
||||
"tests/integration/test_light_automations.py": 18.81,
|
||||
"tests/integration/test_light_binary_effect_off_phase.py": 8.38,
|
||||
"tests/integration/test_light_calls.py": 21.88,
|
||||
"tests/integration/test_light_constant_brightness.py": 59.45,
|
||||
"tests/integration/test_light_control_action.py": 31.91,
|
||||
"tests/integration/test_light_dim_relative_action.py": 14.43,
|
||||
"tests/integration/test_light_effect_zero_brightness.py": 25.05,
|
||||
"tests/integration/test_light_initial_state.py": 18.97,
|
||||
"tests/integration/test_light_toggle_action.py": 17.44,
|
||||
"tests/integration/test_lock_automations.py": 18.9,
|
||||
"tests/integration/test_logger_buffered_recursion_guard.py": 18.2,
|
||||
"tests/integration/test_loop_disable_enable.py": 63.35,
|
||||
"tests/integration/test_loop_interval_decoupling.py": 17.7,
|
||||
"tests/integration/test_loop_interval_default_not_pulled_forward.py": 21.56,
|
||||
"tests/integration/test_micros_to_millis.py": 15.89,
|
||||
"tests/integration/test_multi_click_trigger.py": 17.23,
|
||||
"tests/integration/test_multi_device_preferences.py": 19.4,
|
||||
"tests/integration/test_noise_encryption_key_protection.py": 72.59,
|
||||
"tests/integration/test_object_id_api_verification.py": 19.22,
|
||||
"tests/integration/test_object_id_friendly_name_no_mac_suffix.py": 16.77,
|
||||
"tests/integration/test_object_id_no_friendly_name.py": 45.8,
|
||||
"tests/integration/test_online_image_auto_detects_image_bmp_mime.py": 86.73,
|
||||
"tests/integration/test_online_image_auto_detects_redirected_image_bmp_mime.py": 40.4,
|
||||
"tests/integration/test_online_image_bmp.py": 37.24,
|
||||
"tests/integration/test_oversized_payloads.py": 55.75,
|
||||
"tests/integration/test_preference_key_stability.py": 25.49,
|
||||
"tests/integration/test_runtime_stats.py": 29.81,
|
||||
"tests/integration/test_safe_mode_loop_runs.py": 6.26,
|
||||
"tests/integration/test_scheduler_blocking_warning.py": 37.98,
|
||||
"tests/integration/test_scheduler_bulk_cleanup.py": 18.67,
|
||||
"tests/integration/test_scheduler_defer_cancel.py": 18.46,
|
||||
"tests/integration/test_scheduler_defer_cancel_regular.py": 16.34,
|
||||
"tests/integration/test_scheduler_defer_fifo_simple.py": 18.26,
|
||||
"tests/integration/test_scheduler_defer_stress.py": 17.74,
|
||||
"tests/integration/test_scheduler_heap_stress.py": 3.89,
|
||||
"tests/integration/test_scheduler_internal_id_no_collision.py": 20.01,
|
||||
"tests/integration/test_scheduler_interval_reschedule.py": 16.29,
|
||||
"tests/integration/test_scheduler_interval_zero_coerced.py": 16.09,
|
||||
"tests/integration/test_scheduler_null_name.py": 14.69,
|
||||
"tests/integration/test_scheduler_numeric_id_test.py": 17.08,
|
||||
"tests/integration/test_scheduler_pool.py": 19.88,
|
||||
"tests/integration/test_scheduler_rapid_cancellation.py": 4.42,
|
||||
"tests/integration/test_scheduler_recursive_timeout.py": 4.3,
|
||||
"tests/integration/test_scheduler_removed_item_race.py": 15.49,
|
||||
"tests/integration/test_scheduler_self_keyed.py": 25.77,
|
||||
"tests/integration/test_scheduler_simultaneous_callbacks.py": 14.84,
|
||||
"tests/integration/test_scheduler_string_test.py": 15.42,
|
||||
"tests/integration/test_script_array_params.py": 12.73,
|
||||
"tests/integration/test_script_delay_params.py": 12.69,
|
||||
"tests/integration/test_script_queued.py": 20.38,
|
||||
"tests/integration/test_script_queued_idle_loop.py": 25.06,
|
||||
"tests/integration/test_script_wait_on_boot.py": 15.67,
|
||||
"tests/integration/test_select_stringref_trigger.py": 19.48,
|
||||
"tests/integration/test_sensor_filters_delta.py": 27.62,
|
||||
"tests/integration/test_sensor_filters_ring_buffer.py": 20.27,
|
||||
"tests/integration/test_sensor_filters_sliding_window.py": 56.28,
|
||||
"tests/integration/test_sensor_filters_value_list.py": 20.6,
|
||||
"tests/integration/test_sensor_timeout_filter.py": 22.21,
|
||||
"tests/integration/test_socket_wake_gate_tcp.py": 16.37,
|
||||
"tests/integration/test_status_flags.py": 29.68,
|
||||
"tests/integration/test_strftime_to.py": 17.42,
|
||||
"tests/integration/test_syslog.py": 18.39,
|
||||
"tests/integration/test_template_alarm_control_panel_many_sensors.py": 25.61,
|
||||
"tests/integration/test_template_text_save.py": 19.16,
|
||||
"tests/integration/test_text_command.py": 16.43,
|
||||
"tests/integration/test_text_sensor_raw_state.py": 17.19,
|
||||
"tests/integration/test_uart_mock_ld2410.py": 37.0,
|
||||
"tests/integration/test_uart_mock_ld2412.py": 40.82,
|
||||
"tests/integration/test_uart_mock_ld2420.py": 32.7,
|
||||
"tests/integration/test_uart_mock_ld2450.py": 32.84,
|
||||
"tests/integration/test_uart_mock_modbus.py": 548.87,
|
||||
"tests/integration/test_udp.py": 16.67,
|
||||
"tests/integration/test_use_address_runtime.py": 27.26,
|
||||
"tests/integration/test_valve_control_action.py": 24.58,
|
||||
"tests/integration/test_varint_five_byte_device_id.py": 22.5,
|
||||
"tests/integration/test_wait_until_mid_loop_timing.py": 22.05,
|
||||
"tests/integration/test_wait_until_on_boot.py": 10.37,
|
||||
"tests/integration/test_wait_until_ordering.py": 18.23,
|
||||
"tests/integration/test_wait_until_reentrant_restart.py": 19.35,
|
||||
"tests/integration/test_wake_loop_forces_phase_b.py": 17.83,
|
||||
"tests/integration/test_water_heater_template.py": 25.7
|
||||
"tests/integration/test_action_concurrent_reentry.py": 35.41,
|
||||
"tests/integration/test_addressable_light_transition.py": 27.93,
|
||||
"tests/integration/test_alarm_control_panel_state_transitions.py": 28.03,
|
||||
"tests/integration/test_api_action_metadata.py": 22.44,
|
||||
"tests/integration/test_api_action_responses.py": 23.27,
|
||||
"tests/integration/test_api_action_timeout.py": 42.69,
|
||||
"tests/integration/test_api_conditional_memory.py": 12.77,
|
||||
"tests/integration/test_api_custom_services.py": 17.59,
|
||||
"tests/integration/test_api_get_time_response_timezone.py": 24.75,
|
||||
"tests/integration/test_api_homeassistant.py": 22.48,
|
||||
"tests/integration/test_api_homeassistant_action_no_subscriber.py": 16.41,
|
||||
"tests/integration/test_api_homeassistant_binary_sensor_initial_state.py": 15.96,
|
||||
"tests/integration/test_api_list_entities_backpressure.py": 15.54,
|
||||
"tests/integration/test_api_message_size_batching.py": 19.0,
|
||||
"tests/integration/test_api_reboot_timeout.py": 32.66,
|
||||
"tests/integration/test_api_string_lambda.py": 17.33,
|
||||
"tests/integration/test_api_vv_logging.py": 29.79,
|
||||
"tests/integration/test_api_zero_psk_provisioning.py": 45.19,
|
||||
"tests/integration/test_areas_and_devices.py": 21.24,
|
||||
"tests/integration/test_automation_wait_actions.py": 17.02,
|
||||
"tests/integration/test_automations.py": 40.66,
|
||||
"tests/integration/test_batch_delay_zero_rapid_transitions.py": 26.8,
|
||||
"tests/integration/test_binary_sensor_autorepeat_filter.py": 26.53,
|
||||
"tests/integration/test_binary_sensor_invalidate_state.py": 23.19,
|
||||
"tests/integration/test_blocking_warning_log_time_not_charged_to_next_operation.py": 23.68,
|
||||
"tests/integration/test_build_info.py": 27.63,
|
||||
"tests/integration/test_camera_mock.py": 18.67,
|
||||
"tests/integration/test_climate_control_action.py": 18.7,
|
||||
"tests/integration/test_climate_custom_modes.py": 20.77,
|
||||
"tests/integration/test_continuation_actions.py": 26.26,
|
||||
"tests/integration/test_cover_control_action.py": 26.87,
|
||||
"tests/integration/test_crc8_helper.py": 17.66,
|
||||
"tests/integration/test_device_id_in_state.py": 45.18,
|
||||
"tests/integration/test_duplicate_entities.py": 21.6,
|
||||
"tests/integration/test_entity_icon.py": 35.33,
|
||||
"tests/integration/test_fan_turn_on_action.py": 25.97,
|
||||
"tests/integration/test_fnv1_hash_object_id.py": 18.22,
|
||||
"tests/integration/test_fnv1a_hash.py": 24.44,
|
||||
"tests/integration/test_gpio_expander_cache.py": 21.09,
|
||||
"tests/integration/test_host_logger_thread_safety.py": 17.54,
|
||||
"tests/integration/test_host_mode_basic.py": 3.66,
|
||||
"tests/integration/test_host_mode_batch_delay.py": 25.3,
|
||||
"tests/integration/test_host_mode_climate_basic_state.py": 30.55,
|
||||
"tests/integration/test_host_mode_climate_control.py": 31.32,
|
||||
"tests/integration/test_host_mode_empty_string_options.py": 22.31,
|
||||
"tests/integration/test_host_mode_entity_fields.py": 33.88,
|
||||
"tests/integration/test_host_mode_fan_preset.py": 26.37,
|
||||
"tests/integration/test_host_mode_many_entities.py": 54.44,
|
||||
"tests/integration/test_host_mode_many_entities_multiple_connections.py": 32.68,
|
||||
"tests/integration/test_host_mode_noise_encryption.py": 26.68,
|
||||
"tests/integration/test_host_mode_reconnect.py": 14.53,
|
||||
"tests/integration/test_host_mode_sensor.py": 16.2,
|
||||
"tests/integration/test_host_ota.py": 119.87,
|
||||
"tests/integration/test_host_preferences.py": 27.57,
|
||||
"tests/integration/test_host_preferences_suspend_resume.py": 24.04,
|
||||
"tests/integration/test_improv_serial_uart.py": 41.33,
|
||||
"tests/integration/test_large_message_batching.py": 21.34,
|
||||
"tests/integration/test_legacy_area.py": 14.9,
|
||||
"tests/integration/test_legacy_climate_compat.py": 23.55,
|
||||
"tests/integration/test_legacy_fan_compat.py": 15.08,
|
||||
"tests/integration/test_light_automations.py": 19.32,
|
||||
"tests/integration/test_light_binary_effect_off_phase.py": 25.12,
|
||||
"tests/integration/test_light_calls.py": 35.45,
|
||||
"tests/integration/test_light_constant_brightness.py": 19.95,
|
||||
"tests/integration/test_light_control_action.py": 19.88,
|
||||
"tests/integration/test_light_dim_relative_action.py": 30.87,
|
||||
"tests/integration/test_light_effect_zero_brightness.py": 26.06,
|
||||
"tests/integration/test_light_initial_state.py": 23.99,
|
||||
"tests/integration/test_light_toggle_action.py": 32.41,
|
||||
"tests/integration/test_lock_automations.py": 16.87,
|
||||
"tests/integration/test_logger_buffered_recursion_guard.py": 25.61,
|
||||
"tests/integration/test_loop_disable_enable.py": 23.03,
|
||||
"tests/integration/test_loop_interval_decoupling.py": 18.9,
|
||||
"tests/integration/test_loop_interval_default_not_pulled_forward.py": 26.3,
|
||||
"tests/integration/test_lvgl_headless_render.py": 96.15,
|
||||
"tests/integration/test_micros_to_millis.py": 22.83,
|
||||
"tests/integration/test_multi_click_trigger.py": 27.57,
|
||||
"tests/integration/test_multi_device_preferences.py": 29.63,
|
||||
"tests/integration/test_noise_encryption_key_protection.py": 26.01,
|
||||
"tests/integration/test_object_id_api_verification.py": 20.02,
|
||||
"tests/integration/test_object_id_friendly_name_no_mac_suffix.py": 22.49,
|
||||
"tests/integration/test_object_id_no_friendly_name.py": 42.76,
|
||||
"tests/integration/test_online_image_auto_detects_image_bmp_mime.py": 42.17,
|
||||
"tests/integration/test_online_image_auto_detects_redirected_image_bmp_mime.py": 51.86,
|
||||
"tests/integration/test_online_image_bmp.py": 26.52,
|
||||
"tests/integration/test_oversized_payloads.py": 52.25,
|
||||
"tests/integration/test_preference_key_stability.py": 19.23,
|
||||
"tests/integration/test_runtime_stats.py": 21.0,
|
||||
"tests/integration/test_safe_mode_loop_runs.py": 23.94,
|
||||
"tests/integration/test_scheduler_blocking_warning.py": 32.75,
|
||||
"tests/integration/test_scheduler_bulk_cleanup.py": 23.77,
|
||||
"tests/integration/test_scheduler_defer_cancel.py": 16.93,
|
||||
"tests/integration/test_scheduler_defer_cancel_regular.py": 18.08,
|
||||
"tests/integration/test_scheduler_defer_fifo_simple.py": 18.48,
|
||||
"tests/integration/test_scheduler_defer_stress.py": 19.04,
|
||||
"tests/integration/test_scheduler_heap_stress.py": 28.16,
|
||||
"tests/integration/test_scheduler_internal_id_no_collision.py": 19.65,
|
||||
"tests/integration/test_scheduler_interval_reschedule.py": 21.83,
|
||||
"tests/integration/test_scheduler_interval_zero_coerced.py": 15.55,
|
||||
"tests/integration/test_scheduler_null_name.py": 23.95,
|
||||
"tests/integration/test_scheduler_numeric_id_test.py": 18.59,
|
||||
"tests/integration/test_scheduler_pool.py": 17.91,
|
||||
"tests/integration/test_scheduler_rapid_cancellation.py": 27.58,
|
||||
"tests/integration/test_scheduler_recursive_timeout.py": 15.8,
|
||||
"tests/integration/test_scheduler_removed_item_race.py": 26.03,
|
||||
"tests/integration/test_scheduler_self_keyed.py": 25.84,
|
||||
"tests/integration/test_scheduler_simultaneous_callbacks.py": 16.69,
|
||||
"tests/integration/test_scheduler_string_test.py": 24.02,
|
||||
"tests/integration/test_script_array_params.py": 5.82,
|
||||
"tests/integration/test_script_delay_params.py": 16.29,
|
||||
"tests/integration/test_script_queued.py": 19.29,
|
||||
"tests/integration/test_script_queued_idle_loop.py": 4.6,
|
||||
"tests/integration/test_script_wait_on_boot.py": 17.22,
|
||||
"tests/integration/test_sdl_headless_screenshot.py": 12.92,
|
||||
"tests/integration/test_select_stringref_trigger.py": 25.88,
|
||||
"tests/integration/test_sensor_filters_delta.py": 18.51,
|
||||
"tests/integration/test_sensor_filters_ring_buffer.py": 16.08,
|
||||
"tests/integration/test_sensor_filters_sliding_window.py": 77.89,
|
||||
"tests/integration/test_sensor_filters_value_list.py": 27.78,
|
||||
"tests/integration/test_sensor_timeout_filter.py": 30.05,
|
||||
"tests/integration/test_snapshot_display.py": 26.4,
|
||||
"tests/integration/test_socket_wake_gate_tcp.py": 9.97,
|
||||
"tests/integration/test_status_flags.py": 28.8,
|
||||
"tests/integration/test_strftime_to.py": 25.91,
|
||||
"tests/integration/test_syslog.py": 16.24,
|
||||
"tests/integration/test_template_alarm_control_panel_many_sensors.py": 19.89,
|
||||
"tests/integration/test_template_climate_basic.py": 28.44,
|
||||
"tests/integration/test_template_climate_custom_modes.py": 29.59,
|
||||
"tests/integration/test_template_climate_nonoptimistic.py": 17.75,
|
||||
"tests/integration/test_template_climate_on_control_ordering.py": 16.55,
|
||||
"tests/integration/test_template_climate_publish_all_fields.py": 22.91,
|
||||
"tests/integration/test_template_climate_sensor_push.py": 21.37,
|
||||
"tests/integration/test_template_climate_set_actions.py": 26.3,
|
||||
"tests/integration/test_template_climate_two_point_temperature.py": 19.67,
|
||||
"tests/integration/test_template_text_save.py": 25.04,
|
||||
"tests/integration/test_text_command.py": 24.39,
|
||||
"tests/integration/test_text_sensor_raw_state.py": 18.82,
|
||||
"tests/integration/test_uart_mock_ld2410.py": 46.84,
|
||||
"tests/integration/test_uart_mock_ld2412.py": 80.95,
|
||||
"tests/integration/test_uart_mock_ld2420.py": 43.29,
|
||||
"tests/integration/test_uart_mock_ld2450.py": 16.06,
|
||||
"tests/integration/test_uart_mock_modbus.py": 652.45,
|
||||
"tests/integration/test_udp.py": 9.41,
|
||||
"tests/integration/test_use_address_runtime.py": 27.5,
|
||||
"tests/integration/test_valve_control_action.py": 17.18,
|
||||
"tests/integration/test_varint_five_byte_device_id.py": 25.37,
|
||||
"tests/integration/test_wait_until_mid_loop_timing.py": 14.91,
|
||||
"tests/integration/test_wait_until_on_boot.py": 23.0,
|
||||
"tests/integration/test_wait_until_ordering.py": 14.75,
|
||||
"tests/integration/test_wait_until_reentrant_restart.py": 15.4,
|
||||
"tests/integration/test_wake_loop_forces_phase_b.py": 18.56,
|
||||
"tests/integration/test_water_heater_template.py": 26.2
|
||||
}
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
"""Tests for the espnow component's final validation."""
|
||||
|
||||
import pytest
|
||||
|
||||
from esphome.components.esp32.const import (
|
||||
VARIANT_ESP32C3,
|
||||
VARIANT_ESP32H2,
|
||||
VARIANT_ESP32P4,
|
||||
)
|
||||
from esphome.components.espnow import _validate_variant
|
||||
import esphome.config_validation as cv
|
||||
import esphome.final_validate as fv
|
||||
from esphome.types import ConfigType
|
||||
|
||||
|
||||
def _run(
|
||||
monkeypatch, variant: str, full_config: dict, config: ConfigType
|
||||
) -> ConfigType:
|
||||
monkeypatch.setattr("esphome.components.espnow.get_esp32_variant", lambda: variant)
|
||||
token = fv.full_config.set(full_config)
|
||||
try:
|
||||
return _validate_variant(config)
|
||||
finally:
|
||||
fv.full_config.reset(token)
|
||||
|
||||
|
||||
def test_variant_with_native_wifi_passes(monkeypatch) -> None:
|
||||
"""A variant with a native Wi-Fi PHY needs no shim; config passes through."""
|
||||
config = {"id": "espnow"}
|
||||
assert _run(monkeypatch, VARIANT_ESP32C3, {}, config) is config
|
||||
|
||||
|
||||
def test_radioless_non_p4_variant_rejected(monkeypatch) -> None:
|
||||
"""Radio-less variants without any ESP-NOW path are rejected outright."""
|
||||
with pytest.raises(cv.Invalid, match="not supported"):
|
||||
_run(monkeypatch, VARIANT_ESP32H2, {}, {})
|
||||
|
||||
|
||||
def test_p4_without_esp32_hosted_rejected(monkeypatch) -> None:
|
||||
"""The P4 needs the esp32_hosted shim to supply the esp_now_* symbols."""
|
||||
with pytest.raises(cv.Invalid, match="esp32_hosted"):
|
||||
_run(monkeypatch, VARIANT_ESP32P4, {}, {})
|
||||
|
||||
|
||||
def test_p4_with_esp32_hosted_passes(monkeypatch) -> None:
|
||||
"""The P4 with esp32_hosted present validates; config passes through."""
|
||||
config = {"id": "espnow"}
|
||||
assert _run(monkeypatch, VARIANT_ESP32P4, {"esp32_hosted": {}}, config) is config
|
||||
@@ -9,7 +9,7 @@ not be part of a unit test suite.
|
||||
|
||||
"""
|
||||
|
||||
from collections.abc import Generator
|
||||
from collections.abc import Callable, Generator
|
||||
import os
|
||||
from pathlib import Path
|
||||
import sys
|
||||
@@ -137,3 +137,40 @@ def mock_get_component() -> Generator[Mock, None, None]:
|
||||
"""Mock get_component for config module."""
|
||||
with patch("esphome.config.get_component") as mock:
|
||||
yield mock
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def held_lock() -> Callable[..., Callable[..., None]]:
|
||||
"""Factory for a ``FileLock.acquire`` fake held by another downloader.
|
||||
|
||||
Each poll writes the next chunk to ``part`` (or runs it, for a callable)
|
||||
and raises ``Timeout``; when the chunks run out the part is removed,
|
||||
``land()`` runs, and the acquire succeeds (also for any later job, so
|
||||
``land`` must be idempotent).
|
||||
"""
|
||||
from filelock import Timeout
|
||||
|
||||
def make(
|
||||
part: Path,
|
||||
chunks: list[bytes | Callable[[], None]],
|
||||
land: Callable[[], None],
|
||||
) -> Callable[..., None]:
|
||||
polls = iter(chunks)
|
||||
|
||||
def acquire(*args, **kwargs) -> None:
|
||||
try:
|
||||
chunk = next(polls)
|
||||
except StopIteration:
|
||||
part.unlink(missing_ok=True)
|
||||
land()
|
||||
return
|
||||
if callable(chunk):
|
||||
chunk()
|
||||
else:
|
||||
part.parent.mkdir(parents=True, exist_ok=True)
|
||||
part.write_bytes(chunk)
|
||||
raise Timeout("held")
|
||||
|
||||
return acquire
|
||||
|
||||
return make
|
||||
|
||||
@@ -2353,3 +2353,20 @@ def test_discard_partial_download_logs_undeletable(
|
||||
):
|
||||
framework_helpers.discard_partial_download(dest)
|
||||
assert "Could not remove" in caplog.text
|
||||
|
||||
|
||||
def test_downloaded_bytes_reports_what_is_on_disk(tmp_path: Path) -> None:
|
||||
"""Part file first, then the landed file, both capped at size; else 0."""
|
||||
dest = tmp_path / "archive"
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 0
|
||||
part = tmp_path / "archive.part"
|
||||
part.write_bytes(b"ab")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 2
|
||||
part.write_bytes(b"abcdef")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 4
|
||||
part.unlink()
|
||||
dest.write_bytes(b"abc")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 3
|
||||
assert framework_helpers.downloaded_bytes(dest) == 3
|
||||
dest.write_bytes(b"abcdef")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 4
|
||||
|
||||
@@ -454,23 +454,96 @@ def test_uri_fetch_job_waits_out_a_briefly_held_lock(tmp_path: Path) -> None:
|
||||
assert dl_path.read_bytes() == b"data"
|
||||
|
||||
|
||||
def test_lock_deadline_leaves_download_to_the_holder(tmp_path: Path) -> None:
|
||||
"""A lock held past the deadline means another process is fetching the
|
||||
same file; skipping cleanly beats a misleading failure warning. The
|
||||
tracker is still polled so a parked worker observes cancellation."""
|
||||
@pytest.mark.parametrize("staged", [b"", b"ab"])
|
||||
def test_lock_deadline_leaves_download_to_the_holder(
|
||||
tmp_path: Path, staged: bytes
|
||||
) -> None:
|
||||
"""A lock held past the deadline is another process's download; skip
|
||||
cleanly, polling the tracker with what the holder has staged so far."""
|
||||
dl_path = tmp_path / "archive"
|
||||
(tmp_path / "archive.prefetch.part").write_bytes(staged)
|
||||
ticks: list[int] = []
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch.object(pf, "_DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
):
|
||||
pf._uri_fetch_job(MagicMock(), "https://x/a.zip", dl_path, 4)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == [0]
|
||||
assert ticks == [len(staged)]
|
||||
assert not dl_path.exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("job", "part_name", "chunks", "expected"),
|
||||
[
|
||||
(
|
||||
lambda dl_path: pf._registry_fetch_job(
|
||||
MagicMock(), "https://x/a.tar.gz", dl_path, "ab" * 32, 4
|
||||
),
|
||||
"archive.part",
|
||||
[b"a", b"abc"],
|
||||
[1, 3, 4],
|
||||
),
|
||||
(
|
||||
lambda dl_path: pf._uri_fetch_job(
|
||||
MagicMock(), "https://x/a.zip", dl_path, 4
|
||||
),
|
||||
"archive.prefetch.part",
|
||||
[b"ab"],
|
||||
[2, 4],
|
||||
),
|
||||
],
|
||||
ids=["registry", "uri"],
|
||||
)
|
||||
def test_lock_wait_reports_the_holders_progress(
|
||||
tmp_path: Path,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
held_lock,
|
||||
job,
|
||||
part_name: str,
|
||||
chunks: list[bytes],
|
||||
expected: list[int],
|
||||
) -> None:
|
||||
"""A waiting job reports the holder's part file (the staging one for a
|
||||
URL job), then the full size once the holder lands the archive."""
|
||||
dl_path = tmp_path / "archive"
|
||||
ticks: list[int] = []
|
||||
acquire = held_lock(
|
||||
tmp_path / part_name, chunks, lambda: dl_path.write_bytes(b"abcd")
|
||||
)
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
caplog.at_level(logging.INFO),
|
||||
):
|
||||
job(dl_path)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == expected
|
||||
assert caplog.text.count("Waiting for another process downloading archive") == 1
|
||||
|
||||
|
||||
def test_uri_lock_wait_prefers_the_landed_archive(tmp_path: Path, held_lock) -> None:
|
||||
"""Between the holder's promotion rename and its release the staging
|
||||
part is gone; the landed cache file is credited instead of 0."""
|
||||
dl_path = tmp_path / "archive"
|
||||
ticks: list[int] = []
|
||||
acquire = held_lock(
|
||||
tmp_path / "archive.prefetch.part",
|
||||
[b"ab", lambda: dl_path.write_bytes(b"abcd")],
|
||||
lambda: None,
|
||||
)
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
):
|
||||
pf._uri_fetch_job(MagicMock(), "https://x/a.zip", dl_path, 4)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == [2, 4, 4]
|
||||
|
||||
|
||||
def test_registry_lock_deadline_skips_registration(tmp_path: Path) -> None:
|
||||
"""A registry job that lost the download race to another process
|
||||
must not stamp a nonexistent archive into pio's usage.db."""
|
||||
@@ -479,7 +552,7 @@ def test_registry_lock_deadline_skips_registration(tmp_path: Path) -> None:
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch.object(pf, "_DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
):
|
||||
pf._registry_fetch_job(manager, "https://x/a.tar.gz", dl_path, "ab" * 32, 4)(
|
||||
lambda done: None
|
||||
|
||||
@@ -8,6 +8,7 @@ import os
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from filelock import Timeout
|
||||
import pytest
|
||||
|
||||
from esphome.core import EsphomeError
|
||||
@@ -540,16 +541,13 @@ def test_prefetch_packages_skips_freshly_installed_dest(tmp_path: Path) -> None:
|
||||
dest = tmp_path / "a"
|
||||
dest.mkdir()
|
||||
|
||||
from contextlib import contextmanager
|
||||
|
||||
@contextmanager
|
||||
def marker_appears_under_lock(path, **kwargs):
|
||||
def marker_appears_under_lock(*args, **kwargs):
|
||||
# Simulates the concurrent build finishing while we waited
|
||||
(dest / ".esphome_extracted").touch()
|
||||
yield
|
||||
|
||||
with (
|
||||
patch("filelock.FileLock", side_effect=marker_appears_under_lock),
|
||||
patch("filelock.FileLock.acquire", side_effect=marker_appears_under_lock),
|
||||
patch("filelock.FileLock.release"),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10})
|
||||
@@ -559,6 +557,69 @@ def test_prefetch_packages_skips_freshly_installed_dest(tmp_path: Path) -> None:
|
||||
mock_download.assert_not_called()
|
||||
|
||||
|
||||
def test_prefetch_packages_waits_with_the_holders_progress(
|
||||
tmp_path: Path, held_lock
|
||||
) -> None:
|
||||
"""A worker parked on another build's lock reports that build's part
|
||||
file, then the full size once the marker appears."""
|
||||
dest = tmp_path / "a"
|
||||
dest.mkdir()
|
||||
ticks: list[int] = []
|
||||
part = tmp_path / "dl" / "a-1.0.part"
|
||||
|
||||
def installed_and_pruned() -> None:
|
||||
# install_package touches the marker, then unlinks the archive
|
||||
(dest / ".esphome_extracted").touch()
|
||||
part.unlink()
|
||||
|
||||
acquire = held_lock(
|
||||
part,
|
||||
[lambda: None, b"abc", installed_and_pruned],
|
||||
(dest / ".esphome_extracted").touch,
|
||||
)
|
||||
|
||||
def fake_batch(header, jobs):
|
||||
for _name, _size, fetch in jobs:
|
||||
fetch(ticks.append)
|
||||
return []
|
||||
|
||||
with (
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
patch.object(registry, "run_batch_downloads", side_effect=fake_batch),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10, "b": 5})
|
||||
),
|
||||
):
|
||||
registry.prefetch_packages(
|
||||
[("a", "1.0", dest, []), ("b", "2.0", tmp_path / "b", [])],
|
||||
tmp_path / "dl",
|
||||
)
|
||||
assert ticks == [0, 3, 10, 10]
|
||||
mock_download.assert_called_once()
|
||||
|
||||
|
||||
def test_prefetch_packages_leaves_a_long_held_lock_to_its_holder(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Past the deadline the worker skips; install_package waits on the same
|
||||
lock later and verifies whatever the holder produced."""
|
||||
with (
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10, "b": 5})
|
||||
),
|
||||
):
|
||||
registry.prefetch_packages(
|
||||
[("a", "1.0", tmp_path / "a", []), ("b", "2.0", tmp_path / "b", [])],
|
||||
tmp_path / "dl",
|
||||
)
|
||||
mock_download.assert_not_called()
|
||||
|
||||
|
||||
def test_already_installed_probe(tmp_path: Path) -> None:
|
||||
"""Both arms of the marker probe the prefetch worker keys on."""
|
||||
dest = tmp_path / "pkg"
|
||||
|
||||
Reference in New Issue
Block a user